mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +00:00
docs(plugin-arch): 归档插件架构迁移评估 + plan 第11节整改计划
- docs/zh/架构迁移评估.md: C ABI→子进程+共享内存完整迁移论证(1621行) - docs/zh/experiments/: 18项可复跑可行性实验(架构评估的所有数字来源) - plan.md §11: 11.1~11.9 插件架构缺陷修复清单(唯一权威编号) - main 保持干净,本批次为 update 特性分支的整改起点
This commit is contained in:
49
docs/zh/experiments/plugin-arch/02-feasibility/exp10.go
Normal file
49
docs/zh/experiments/plugin-arch/02-feasibility/exp10.go
Normal file
@ -0,0 +1,49 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"golang.org/x/sys/unix"
|
||||
)
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 10:多媒体 payload —— 共享内存零拷贝 vs JSON base64 ===")
|
||||
sizes := []int{100 * 1024, 1024 * 1024, 5 * 1024 * 1024}
|
||||
for _, sz := range sizes {
|
||||
img := make([]byte, sz)
|
||||
for i := range img { img[i] = byte(i % 251) }
|
||||
|
||||
// A. JSON + base64(当前 ContentBlock 的做法)
|
||||
t0 := time.Now()
|
||||
b64 := base64.StdEncoding.EncodeToString(img)
|
||||
blob, _ := json.Marshal(map[string]string{"type": "image_url", "url": "data:image/png;base64," + b64})
|
||||
var back map[string]string
|
||||
json.Unmarshal(blob, &back)
|
||||
dec, _ := base64.StdEncoding.DecodeString(back["url"][22:])
|
||||
jsonDur := time.Since(t0)
|
||||
|
||||
// B. 共享内存 arena(写入 + 偏移解引用,零拷贝读)
|
||||
mfd, _ := unix.MemfdCreate("arena", 0)
|
||||
unix.Ftruncate(mfd, int64(sz+4096))
|
||||
data, _ := unix.Mmap(mfd, 0, sz+4096, unix.PROT_READ|unix.PROT_WRITE, unix.MAP_SHARED)
|
||||
t0 = time.Now()
|
||||
copy(data[4096:], img) // 写 arena
|
||||
view := data[4096 : 4096+sz] // 偏移解引用 = 零拷贝切片
|
||||
_ = view[sz-1]
|
||||
shmDur := time.Since(t0)
|
||||
unix.Munmap(data)
|
||||
unix.Close(mfd)
|
||||
|
||||
fmt.Printf("\n%s payload:\n", map[int]string{100*1024:"100KB", 1024*1024:"1MB", 5*1024*1024:"5MB"}[sz])
|
||||
fmt.Printf(" A JSON+base64: %8v 传输体积 %d B (+%.0f%%) 解出 %d B %s\n",
|
||||
jsonDur, len(blob), float64(len(blob)-sz)/float64(sz)*100, len(dec),
|
||||
map[bool]string{true:"✓",false:"✗"}[len(dec)==sz])
|
||||
fmt.Printf(" B 共享内存: %8v 传输体积 8 B (描述符) 零拷贝视图 %d B\n", shmDur, len(view))
|
||||
fmt.Printf(" → 加速 %.0fx, 体积节省 %.0f%%\n",
|
||||
float64(jsonDur)/float64(shmDur), float64(len(blob)-8)/float64(len(blob))*100)
|
||||
}
|
||||
}
|
||||
42
docs/zh/experiments/plugin-arch/02-feasibility/exp11.go
Normal file
42
docs/zh/experiments/plugin-arch/02-feasibility/exp11.go
Normal file
@ -0,0 +1,42 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os/exec"
|
||||
"sort"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Req struct{ ID int `json:"id"`; Method string `json:"method"`; Args json.RawMessage `json:"args"` }
|
||||
type Res struct{ ID int `json:"id"`; Result string `json:"result"` }
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 11:工具调用 RPC 端到端延迟(实测 payload 中位 93B)===")
|
||||
cmd := exec.Command("./plug11")
|
||||
sin, _ := cmd.StdinPipe(); sout, _ := cmd.StdoutPipe()
|
||||
cmd.Start()
|
||||
enc := json.NewEncoder(bufio.NewWriter(sin))
|
||||
w := bufio.NewWriter(sin); enc = json.NewEncoder(w)
|
||||
dec := json.NewDecoder(bufio.NewReader(sout))
|
||||
|
||||
args := json.RawMessage(`{"city":"hangzhou","days":3,"unit":"celsius","detail":true}`)
|
||||
const N = 10000
|
||||
lat := make([]time.Duration, 0, N)
|
||||
for i := 0; i < N; i++ {
|
||||
t0 := time.Now()
|
||||
enc.Encode(Req{ID: i, Method: "weather_query", Args: args}); w.Flush()
|
||||
var r Res
|
||||
if err := dec.Decode(&r); err != nil { break }
|
||||
lat = append(lat, time.Since(t0))
|
||||
}
|
||||
sin.Close(); cmd.Wait()
|
||||
sort.Slice(lat, func(a,b int) bool { return lat[a] < lat[b] })
|
||||
p := func(q float64) time.Duration { return lat[int(float64(len(lat))*q)] }
|
||||
fmt.Printf("样本 %d 次\n", len(lat))
|
||||
fmt.Printf(" p50 = %v\n p90 = %v\n p99 = %v\n max = %v\n", p(0.5), p(0.9), p(0.99), lat[len(lat)-1])
|
||||
fmt.Printf("\n对照 LLM 单轮往返 2-8 秒 → RPC 占比 ≈ %.5f%%\n",
|
||||
float64(p(0.5))/float64(3*time.Second)*100)
|
||||
}
|
||||
12
docs/zh/experiments/plugin-arch/02-feasibility/exp11_plug.go
Normal file
12
docs/zh/experiments/plugin-arch/02-feasibility/exp11_plug.go
Normal file
@ -0,0 +1,12 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
import ("bufio";"encoding/json";"os")
|
||||
type Req struct{ ID int `json:"id"`; Method string `json:"method"`; Args json.RawMessage `json:"args"` }
|
||||
type Res struct{ ID int `json:"id"`; Result string `json:"result"` }
|
||||
func main(){
|
||||
dec:=json.NewDecoder(bufio.NewReader(os.Stdin))
|
||||
w:=bufio.NewWriter(os.Stdout); enc:=json.NewEncoder(w)
|
||||
for { var q Req
|
||||
if err:=dec.Decode(&q); err!=nil {return}
|
||||
enc.Encode(Res{ID:q.ID, Result:`{"ok":true,"data":"` + string(q.Args) + `"}`}); w.Flush() }
|
||||
}
|
||||
@ -0,0 +1,58 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"runtime"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"golang.org/x/sys/unix"
|
||||
)
|
||||
|
||||
func threads() int { e, _ := os.ReadDir("/proc/self/task"); return len(e) }
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 1:eventfd 是否走 Go netpoller(只 park goroutine 不占 OS 线程)===")
|
||||
base := threads()
|
||||
fmt.Printf("基线线程数: %d (GOMAXPROCS=%d)\n\n", base, runtime.GOMAXPROCS(0))
|
||||
|
||||
const N = 200 // 模拟 200 个订阅者等待
|
||||
var wg sync.WaitGroup
|
||||
var woke int64
|
||||
files := make([]*os.File, N)
|
||||
|
||||
for i := 0; i < N; i++ {
|
||||
efd, err := unix.Eventfd(0, unix.EFD_NONBLOCK|unix.EFD_CLOEXEC)
|
||||
if err != nil { fmt.Println("eventfd 失败:", err); return }
|
||||
f := os.NewFile(uintptr(efd), fmt.Sprintf("evt%d", i))
|
||||
files[i] = f
|
||||
wg.Add(1)
|
||||
go func(f *os.File) {
|
||||
defer wg.Done()
|
||||
buf := make([]byte, 8)
|
||||
// 阻塞读:若走 netpoller 只 park goroutine
|
||||
if _, err := f.Read(buf); err == nil {
|
||||
atomic.AddInt64(&woke, 1)
|
||||
}
|
||||
}(f)
|
||||
}
|
||||
|
||||
time.Sleep(500 * time.Millisecond) // 让所有 goroutine 进入等待
|
||||
waiting := threads()
|
||||
fmt.Printf("%d 个 goroutine 阻塞在 eventfd.Read 后:\n", N)
|
||||
fmt.Printf(" 线程数 = %d (增长 %d)\n", waiting, waiting-base)
|
||||
if waiting-base < 20 {
|
||||
fmt.Println(" ✅ 走 netpoller:线程未随等待者数量增长")
|
||||
} else {
|
||||
fmt.Printf(" ❌ 退化为阻塞 syscall:每个等待者占一个 OS 线程\n")
|
||||
}
|
||||
|
||||
// 全部唤醒
|
||||
one := []byte{1,0,0,0,0,0,0,0}
|
||||
for _, f := range files { f.Write(one) }
|
||||
wg.Wait()
|
||||
fmt.Printf("\n唤醒数 = %d/%d 唤醒后线程数 = %d\n", woke, N, threads())
|
||||
}
|
||||
41
docs/zh/experiments/plugin-arch/02-feasibility/exp2_child.go
Normal file
41
docs/zh/experiments/plugin-arch/02-feasibility/exp2_child.go
Normal file
@ -0,0 +1,41 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"os"
|
||||
"unsafe"
|
||||
|
||||
"golang.org/x/sys/unix"
|
||||
)
|
||||
|
||||
// 子进程:fd 3 = eventfd(通知), fd 4 = shm 文件
|
||||
func main() {
|
||||
efd := os.NewFile(3, "evt")
|
||||
shmf := os.NewFile(4, "shm")
|
||||
|
||||
data, err := unix.Mmap(int(shmf.Fd()), 0, 4096, unix.PROT_READ|unix.PROT_WRITE, unix.MAP_SHARED)
|
||||
if err != nil { fmt.Println("CHILD mmap 失败:", err); os.Exit(1) }
|
||||
fmt.Printf("CHILD: mmap 基址 = %p\n", unsafe.Pointer(&data[0]))
|
||||
|
||||
buf := make([]byte, 8)
|
||||
if _, err := efd.Read(buf); err != nil {
|
||||
fmt.Println("CHILD read err:", err); os.Exit(1)
|
||||
}
|
||||
n := binary.LittleEndian.Uint64(buf)
|
||||
fmt.Printf("CHILD: 被 eventfd 唤醒, 计数=%d\n", n)
|
||||
|
||||
// 按偏移读:头部 16 字节 = {off uint32, len uint32, seq uint64}
|
||||
off := binary.LittleEndian.Uint32(data[0:4])
|
||||
ln := binary.LittleEndian.Uint32(data[4:8])
|
||||
seq := binary.LittleEndian.Uint64(data[8:16])
|
||||
payload := string(data[off : off+ln])
|
||||
fmt.Printf("CHILD: 偏移解引用 off=%d len=%d seq=%d → %q\n", off, ln, seq, payload)
|
||||
|
||||
// 子进程回写(验证双向可见)
|
||||
copy(data[2048:], []byte("CHILD-ACK"))
|
||||
binary.LittleEndian.PutUint32(data[16:20], 2048)
|
||||
binary.LittleEndian.PutUint32(data[20:24], uint32(len("CHILD-ACK")))
|
||||
fmt.Println("CHILD: 已回写 ACK")
|
||||
}
|
||||
@ -0,0 +1,60 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"time"
|
||||
"unsafe"
|
||||
|
||||
"golang.org/x/sys/unix"
|
||||
)
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 2:跨进程 eventfd 通知 + 共享内存偏移解引用 ===")
|
||||
|
||||
// eventfd 不带 CLOEXEC(需要被子进程继承)
|
||||
efd, err := unix.Eventfd(0, unix.EFD_NONBLOCK)
|
||||
if err != nil { panic(err) }
|
||||
evtFile := os.NewFile(uintptr(efd), "evt")
|
||||
|
||||
// shm: 用 memfd(匿名,无需 /dev/shm 清理)
|
||||
mfd, err := unix.MemfdCreate("stagectx", 0)
|
||||
if err != nil { panic(err) }
|
||||
if err := unix.Ftruncate(mfd, 4096); err != nil { panic(err) }
|
||||
shmFile := os.NewFile(uintptr(mfd), "shm")
|
||||
|
||||
data, err := unix.Mmap(mfd, 0, 4096, unix.PROT_READ|unix.PROT_WRITE, unix.MAP_SHARED)
|
||||
if err != nil { panic(err) }
|
||||
fmt.Printf("PARENT: mmap 基址 = %p\n", unsafe.Pointer(&data[0]))
|
||||
|
||||
// 写 payload 到 arena(偏移 1024),头部记描述符
|
||||
msg := "hello-from-parent-via-offset"
|
||||
copy(data[1024:], []byte(msg))
|
||||
binary.LittleEndian.PutUint32(data[0:4], 1024)
|
||||
binary.LittleEndian.PutUint32(data[4:8], uint32(len(msg)))
|
||||
binary.LittleEndian.PutUint64(data[8:16], 42)
|
||||
fmt.Printf("PARENT: 数据已落地 arena@1024, 描述符 {off:1024, len:%d, seq:42}\n", len(msg))
|
||||
|
||||
cmd := exec.Command("go", "run", "exp2_child.go")
|
||||
cmd.ExtraFiles = []*os.File{evtFile, shmFile} // → 子进程 fd 3, 4
|
||||
cmd.Stdout, cmd.Stderr = os.Stdout, os.Stderr
|
||||
if err := cmd.Start(); err != nil { panic(err) }
|
||||
|
||||
time.Sleep(3 * time.Second) // 等 go run 编译+启动
|
||||
fmt.Println("PARENT: 数据到位后 post eventfd(不等待消费者)")
|
||||
t0 := time.Now()
|
||||
evtFile.Write([]byte{1,0,0,0,0,0,0,0})
|
||||
fmt.Printf("PARENT: post 耗时 %v ← post-and-forget\n", time.Since(t0))
|
||||
|
||||
cmd.Wait()
|
||||
|
||||
// 读子进程回写
|
||||
off := binary.LittleEndian.Uint32(data[16:20])
|
||||
ln := binary.LittleEndian.Uint32(data[20:24])
|
||||
if ln > 0 {
|
||||
fmt.Printf("PARENT: 读到子进程回写 → %q ✅ 双向可见\n", string(data[off:off+ln]))
|
||||
}
|
||||
}
|
||||
31
docs/zh/experiments/plugin-arch/02-feasibility/exp3_child.go
Normal file
31
docs/zh/experiments/plugin-arch/02-feasibility/exp3_child.go
Normal file
@ -0,0 +1,31 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
type req struct{ ID int `json:"id"`; Method string `json:"method"` }
|
||||
type resp struct{ ID int `json:"id"`; OK bool `json:"ok"` }
|
||||
|
||||
func main() {
|
||||
in := bufio.NewReader(os.Stdin)
|
||||
out := bufio.NewWriter(os.Stdout)
|
||||
enc, dec := json.NewEncoder(out), json.NewDecoder(in)
|
||||
|
||||
const N = 20000
|
||||
t0 := time.Now()
|
||||
for i := 0; i < N; i++ {
|
||||
enc.Encode(req{ID: i, Method: "stage.lock"})
|
||||
out.Flush()
|
||||
var r resp
|
||||
if err := dec.Decode(&r); err != nil { fmt.Fprintln(os.Stderr, "dec:", err); return }
|
||||
}
|
||||
d := time.Since(t0)
|
||||
fmt.Fprintf(os.Stderr, "CHILD: %d 次 lock RPC 往返 用时 %v, 均摊 %.2f µs/次\n",
|
||||
N, d, float64(d.Microseconds())/float64(N))
|
||||
}
|
||||
@ -0,0 +1,37 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"sync"
|
||||
)
|
||||
|
||||
type req struct{ ID int `json:"id"`; Method string `json:"method"` }
|
||||
type resp struct{ ID int `json:"id"`; OK bool `json:"ok"` }
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 3:锁仲裁 RPC 往返成本(stdio JSON-RPC)===")
|
||||
cmd := exec.Command("go", "run", "exp3_child.go")
|
||||
stdin, _ := cmd.StdinPipe()
|
||||
stdout, _ := cmd.StdoutPipe()
|
||||
cmd.Stderr = os.Stderr
|
||||
cmd.Start()
|
||||
|
||||
var mu sync.Mutex // 内核侧真实的锁仲裁
|
||||
dec := json.NewDecoder(bufio.NewReader(stdout))
|
||||
w := bufio.NewWriter(stdin)
|
||||
enc := json.NewEncoder(w)
|
||||
for {
|
||||
var q req
|
||||
if err := dec.Decode(&q); err != nil { break }
|
||||
mu.Lock() // 真实加锁
|
||||
mu.Unlock() // 立即释放(模拟仲裁开销)
|
||||
enc.Encode(resp{ID: q.ID, OK: true})
|
||||
w.Flush()
|
||||
}
|
||||
cmd.Wait()
|
||||
}
|
||||
71
docs/zh/experiments/plugin-arch/02-feasibility/exp4.go
Normal file
71
docs/zh/experiments/plugin-arch/02-feasibility/exp4.go
Normal file
@ -0,0 +1,71 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"golang.org/x/sys/unix"
|
||||
)
|
||||
|
||||
type ring struct {
|
||||
writeSeq atomic.Uint64
|
||||
cap uint64
|
||||
slots []uint64
|
||||
}
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 4:事件环 post-and-forget vs 同步 Publish(慢消费者场景)===")
|
||||
const tokens = 5000
|
||||
|
||||
// --- A. 现状:同步 Publish,消费者慢 ---
|
||||
slowHandler := func() { time.Sleep(20 * time.Microsecond) }
|
||||
t0 := time.Now()
|
||||
for i := 0; i < tokens; i++ { slowHandler() }
|
||||
syncDur := time.Since(t0)
|
||||
fmt.Printf("A 同步 Publish (慢消费者 20µs): %d token 耗时 %v → 均摊 %.1f µs/token\n",
|
||||
tokens, syncDur, float64(syncDur.Microseconds())/tokens)
|
||||
|
||||
// --- B. 新方案:写环 + eventfd post,不等消费者 ---
|
||||
r := &ring{cap: 1024, slots: make([]uint64, 1024)}
|
||||
efd, _ := unix.Eventfd(0, unix.EFD_NONBLOCK)
|
||||
f := os.NewFile(uintptr(efd), "e")
|
||||
|
||||
var dropped atomic.Uint64
|
||||
// 慢消费者 goroutine
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
buf := make([]byte, 8)
|
||||
var readSeq uint64
|
||||
for {
|
||||
if _, err := f.Read(buf); err != nil { return }
|
||||
w := r.writeSeq.Load()
|
||||
if w-readSeq > r.cap {
|
||||
dropped.Add(w - readSeq - r.cap)
|
||||
readSeq = w - r.cap
|
||||
}
|
||||
for readSeq < w { readSeq++ }
|
||||
time.Sleep(20 * time.Microsecond) // 慢
|
||||
select { case <-done: return; default: }
|
||||
}
|
||||
}()
|
||||
|
||||
t0 = time.Now()
|
||||
one := []byte{1,0,0,0,0,0,0,0}
|
||||
for i := 0; i < tokens; i++ {
|
||||
s := r.writeSeq.Add(1)
|
||||
r.slots[s%r.cap] = s // 写数据
|
||||
f.Write(one) // post,不等
|
||||
}
|
||||
asyncDur := time.Since(t0)
|
||||
close(done)
|
||||
fmt.Printf("B 环+eventfd post: %d token 耗时 %v → 均摊 %.2f µs/token\n",
|
||||
tokens, asyncDur, float64(asyncDur.Microseconds())/tokens)
|
||||
fmt.Printf("\n加速比 %.1fx 丢弃事件 %d(消费者跟不上,已计数)\n",
|
||||
float64(syncDur)/float64(asyncDur), dropped.Load())
|
||||
if asyncDur < syncDur/5 {
|
||||
fmt.Println("✅ post-and-forget 使流式发布与消费者速度解耦")
|
||||
}
|
||||
}
|
||||
55
docs/zh/experiments/plugin-arch/02-feasibility/exp5.go
Normal file
55
docs/zh/experiments/plugin-arch/02-feasibility/exp5.go
Normal file
@ -0,0 +1,55 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
func pssKB(pid int) int {
|
||||
b, err := os.ReadFile(fmt.Sprintf("/proc/%d/smaps_rollup", pid))
|
||||
if err != nil { return 0 }
|
||||
for _, l := range strings.Split(string(b), "\n") {
|
||||
if strings.HasPrefix(l, "Pss:") {
|
||||
f := strings.Fields(l)
|
||||
n, _ := strconv.Atoi(f[1]); return n
|
||||
}
|
||||
}
|
||||
return 0
|
||||
}
|
||||
func threads(pid int) int {
|
||||
e, _ := os.ReadDir(fmt.Sprintf("/proc/%d/task", pid)); return len(e)
|
||||
}
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 5:17 个 Go 子进程插件的真实常驻开销(PSS 计入共享页去重)===")
|
||||
var cmds []*exec.Cmd
|
||||
for i := 0; i < 17; i++ {
|
||||
c := exec.Command("./plugbin")
|
||||
c.Stdin, _ = os.Open(os.DevNull)
|
||||
if err := c.Start(); err != nil { fmt.Println("start:", err); return }
|
||||
cmds = append(cmds, c)
|
||||
}
|
||||
time.Sleep(1500 * time.Millisecond)
|
||||
|
||||
totalPss, totalThreads := 0, 0
|
||||
for _, c := range cmds {
|
||||
totalPss += pssKB(c.Process.Pid)
|
||||
totalThreads += threads(c.Process.Pid)
|
||||
}
|
||||
fmt.Printf("17 进程合计: PSS = %.1f MB, 线程 = %d\n", float64(totalPss)/1024, totalThreads)
|
||||
fmt.Printf("单进程均摊: PSS = %.2f MB, 线程 = %.1f\n",
|
||||
float64(totalPss)/1024/17, float64(totalThreads)/17)
|
||||
fmt.Printf("\n对照 homed 当前(单进程装 17 个 .so):\n")
|
||||
// 找 homed
|
||||
out, _ := exec.Command("pgrep", "-x", "homed").Output()
|
||||
if p := strings.TrimSpace(string(out)); p != "" {
|
||||
pid, _ := strconv.Atoi(strings.Fields(p)[0])
|
||||
fmt.Printf(" homed PSS = %.1f MB, 线程 = %d\n", float64(pssKB(pid))/1024, threads(pid))
|
||||
}
|
||||
for _, c := range cmds { c.Process.Kill(); c.Wait() }
|
||||
}
|
||||
@ -0,0 +1,23 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"os"
|
||||
)
|
||||
|
||||
// 模拟一个最小插件:stdio JSON-RPC loop + 一个 goroutine
|
||||
func main() {
|
||||
go func() { select {} }()
|
||||
in := bufio.NewReader(os.Stdin)
|
||||
dec := json.NewDecoder(in)
|
||||
out := bufio.NewWriter(os.Stdout)
|
||||
enc := json.NewEncoder(out)
|
||||
for {
|
||||
var m map[string]interface{}
|
||||
if err := dec.Decode(&m); err != nil { return }
|
||||
enc.Encode(map[string]interface{}{"ok": true})
|
||||
out.Flush()
|
||||
}
|
||||
}
|
||||
68
docs/zh/experiments/plugin-arch/02-feasibility/exp5b.go
Normal file
68
docs/zh/experiments/plugin-arch/02-feasibility/exp5b.go
Normal file
@ -0,0 +1,68 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
func pssKB(pid int) int {
|
||||
b, err := os.ReadFile(fmt.Sprintf("/proc/%d/smaps_rollup", pid))
|
||||
if err != nil { return -1 }
|
||||
for _, l := range strings.Split(string(b), "\n") {
|
||||
if strings.HasPrefix(l, "Pss:") { f := strings.Fields(l); n,_ := strconv.Atoi(f[1]); return n }
|
||||
}
|
||||
return -1
|
||||
}
|
||||
func rssKB(pid int) int {
|
||||
b, err := os.ReadFile(fmt.Sprintf("/proc/%d/status", pid))
|
||||
if err != nil { return -1 }
|
||||
for _, l := range strings.Split(string(b), "\n") {
|
||||
if strings.HasPrefix(l, "VmRSS:") { f := strings.Fields(l); n,_ := strconv.Atoi(f[1]); return n }
|
||||
}
|
||||
return -1
|
||||
}
|
||||
func threads(pid int) int { e,_ := os.ReadDir(fmt.Sprintf("/proc/%d/task", pid)); return len(e) }
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 5b:17 个 Go 子进程常驻开销(保持 stdin 管道存活)===")
|
||||
var cmds []*exec.Cmd
|
||||
var pipes []interface{ Close() error }
|
||||
for i := 0; i < 17; i++ {
|
||||
c := exec.Command("./plugbin")
|
||||
w, _ := c.StdinPipe() // 保持打开 → 不 EOF
|
||||
pipes = append(pipes, w)
|
||||
c.Stdout = nil
|
||||
if err := c.Start(); err != nil { fmt.Println(err); return }
|
||||
cmds = append(cmds, c)
|
||||
}
|
||||
time.Sleep(2 * time.Second)
|
||||
|
||||
tp, tr, tt, alive := 0, 0, 0, 0
|
||||
for _, c := range cmds {
|
||||
pid := c.Process.Pid
|
||||
if _, err := os.Stat(fmt.Sprintf("/proc/%d", pid)); err != nil { continue }
|
||||
alive++
|
||||
if v := pssKB(pid); v > 0 { tp += v }
|
||||
if v := rssKB(pid); v > 0 { tr += v }
|
||||
tt += threads(pid)
|
||||
}
|
||||
fmt.Printf("存活进程 %d/17\n", alive)
|
||||
fmt.Printf("合计: PSS=%.1f MB RSS=%.1f MB 线程=%d\n",
|
||||
float64(tp)/1024, float64(tr)/1024, tt)
|
||||
if alive > 0 {
|
||||
fmt.Printf("均摊: PSS=%.2f MB RSS=%.2f MB 线程=%.1f\n",
|
||||
float64(tp)/1024/float64(alive), float64(tr)/1024/float64(alive), float64(tt)/float64(alive))
|
||||
}
|
||||
out, _ := exec.Command("pgrep", "-x", "homed").Output()
|
||||
if p := strings.TrimSpace(string(out)); p != "" {
|
||||
pid, _ := strconv.Atoi(strings.Fields(p)[0])
|
||||
fmt.Printf("\n对照 homed(单进程 + 17 个 .so): RSS=%.1f MB 线程=%d\n",
|
||||
float64(rssKB(pid))/1024, threads(pid))
|
||||
}
|
||||
for _, c := range cmds { c.Process.Kill(); c.Wait() }
|
||||
}
|
||||
49
docs/zh/experiments/plugin-arch/02-feasibility/exp6.go
Normal file
49
docs/zh/experiments/plugin-arch/02-feasibility/exp6.go
Normal file
@ -0,0 +1,49 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"time"
|
||||
)
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 6:子进程崩溃隔离 + 退出码/EOF 作为 recordCrash 信号 ===")
|
||||
cmd := exec.Command("./crashbin")
|
||||
sin, _ := cmd.StdinPipe()
|
||||
sout, _ := cmd.StdoutPipe()
|
||||
cmd.Stderr = nil // 丢弃 panic 栈
|
||||
cmd.Start()
|
||||
fmt.Printf("插件进程 pid=%d 已启动\n", cmd.Process.Pid)
|
||||
|
||||
enc := json.NewEncoder(sin)
|
||||
dec := json.NewDecoder(bufio.NewReader(sout))
|
||||
|
||||
// 正常调用
|
||||
enc.Encode(map[string]string{"method": "ping"})
|
||||
var r map[string]interface{}
|
||||
if err := dec.Decode(&r); err == nil { fmt.Println("正常调用 → ", r) }
|
||||
|
||||
// 触发崩溃
|
||||
fmt.Println("\n发送 boom(插件内 panic)...")
|
||||
t0 := time.Now()
|
||||
enc.Encode(map[string]string{"method": "boom"})
|
||||
err := dec.Decode(&r)
|
||||
|
||||
detected := "未检测到"
|
||||
if errors.Is(err, io.EOF) || err == io.ErrUnexpectedEOF { detected = "EOF" } else if err != nil { detected = fmt.Sprintf("%v", err) }
|
||||
fmt.Printf("调用侧感知: %s (耗时 %v)\n", detected, time.Since(t0))
|
||||
|
||||
werr := cmd.Wait()
|
||||
var ec int = -1
|
||||
if ee, ok := werr.(*exec.ExitError); ok { ec = ee.ExitCode() }
|
||||
fmt.Printf("进程退出码 = %d (panic → 2,可直接喂 recordCrash)\n", ec)
|
||||
|
||||
fmt.Printf("\n宿主进程仍存活: pid=%d ✅ 崩溃已隔离\n", os.Getpid())
|
||||
fmt.Println("→ 对照:当前 .so 模型下,bridge 兜不住的 panic 会带崩整个 homed")
|
||||
}
|
||||
23
docs/zh/experiments/plugin-arch/02-feasibility/exp6_crash.go
Normal file
23
docs/zh/experiments/plugin-arch/02-feasibility/exp6_crash.go
Normal file
@ -0,0 +1,23 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"os"
|
||||
)
|
||||
|
||||
func main() {
|
||||
dec := json.NewDecoder(bufio.NewReader(os.Stdin))
|
||||
out := bufio.NewWriter(os.Stdout)
|
||||
enc := json.NewEncoder(out)
|
||||
for {
|
||||
var m map[string]interface{}
|
||||
if err := dec.Decode(&m); err != nil { return }
|
||||
if m["method"] == "boom" {
|
||||
panic("插件故意崩溃") // 真 panic
|
||||
}
|
||||
enc.Encode(map[string]interface{}{"ok": true})
|
||||
out.Flush()
|
||||
}
|
||||
}
|
||||
63
docs/zh/experiments/plugin-arch/02-feasibility/exp7.go
Normal file
63
docs/zh/experiments/plugin-arch/02-feasibility/exp7.go
Normal file
@ -0,0 +1,63 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"time"
|
||||
)
|
||||
|
||||
func spawnAndAsk(bin string) string {
|
||||
cmd := exec.Command(bin)
|
||||
sin, _ := cmd.StdinPipe()
|
||||
sout, _ := cmd.StdoutPipe()
|
||||
cmd.Start()
|
||||
enc := json.NewEncoder(sin)
|
||||
dec := json.NewDecoder(bufio.NewReader(sout))
|
||||
enc.Encode(map[string]string{"method": "version"})
|
||||
var r map[string]interface{}
|
||||
dec.Decode(&r)
|
||||
sin.Close()
|
||||
cmd.Process.Kill()
|
||||
cmd.Wait()
|
||||
if v, ok := r["version"].(string); ok { return v }
|
||||
return "?"
|
||||
}
|
||||
|
||||
func build(ver, out string) {
|
||||
src := fmt.Sprintf(`package main
|
||||
import ("bufio";"encoding/json";"os")
|
||||
func main(){
|
||||
dec:=json.NewDecoder(bufio.NewReader(os.Stdin))
|
||||
w:=bufio.NewWriter(os.Stdout); enc:=json.NewEncoder(w)
|
||||
for { var m map[string]interface{}
|
||||
if err:=dec.Decode(&m); err!=nil {return}
|
||||
enc.Encode(map[string]string{"version":%q}); w.Flush() }
|
||||
}`, ver)
|
||||
os.MkdirAll("v", 0755)
|
||||
os.WriteFile("v/main.go", []byte(src), 0644)
|
||||
os.WriteFile("v/go.mod", []byte("module v\ngo 1.21\n"), 0644)
|
||||
c := exec.Command("go", "build", "-o", "../"+out, ".")
|
||||
c.Dir = "v"
|
||||
if b, err := c.CombinedOutput(); err != nil { fmt.Println("build err:", string(b)) }
|
||||
}
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 7:子进程模型下的热重载(迁移的原始目标)===")
|
||||
build("v1.0.0", "hotbin")
|
||||
fmt.Printf("1) 首次启动插件 → version = %s\n", spawnAndAsk("./hotbin"))
|
||||
|
||||
fmt.Println("2) 替换二进制为 v2.0.0(同路径,无需版本化 hash 目录)")
|
||||
build("v2.0.0", "hotbin")
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
|
||||
v := spawnAndAsk("./hotbin")
|
||||
fmt.Printf("3) 重启插件进程 → version = %s\n", v)
|
||||
if v == "v2.0.0" {
|
||||
fmt.Println("\n✅ 同路径替换即生效:无 NODELETE、无版本化路径、无线程泄漏")
|
||||
fmt.Println(" 对照 .so 模型:同路径 dlopen 复用旧映像,永远拿不到 v2")
|
||||
}
|
||||
}
|
||||
84
docs/zh/experiments/plugin-arch/02-feasibility/exp8.go
Normal file
84
docs/zh/experiments/plugin-arch/02-feasibility/exp8.go
Normal file
@ -0,0 +1,84 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"golang.org/x/sys/unix"
|
||||
)
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 8:跨进程并发扇出改写同一 StageContext(最高风险点 3.4)===")
|
||||
|
||||
mfd, _ := unix.MemfdCreate("stagectx", 0)
|
||||
unix.Ftruncate(mfd, 65536)
|
||||
shmFile := os.NewFile(uintptr(mfd), "shm")
|
||||
data, _ := unix.Mmap(mfd, 0, 65536, unix.PROT_READ|unix.PROT_WRITE, unix.MAP_SHARED)
|
||||
|
||||
// 初始 final_text = "" @1024, arena 游标 = 1024
|
||||
binary.LittleEndian.PutUint32(data[0:4], 1024)
|
||||
binary.LittleEndian.PutUint32(data[4:8], 0)
|
||||
binary.LittleEndian.PutUint32(data[8:12], 1024)
|
||||
|
||||
tags := []string{"A", "B", "C", "D", "E"} // 5 个并发插件
|
||||
var mu sync.Mutex // 内核侧锁仲裁
|
||||
var wg sync.WaitGroup
|
||||
var rpcCount int64
|
||||
var cntMu sync.Mutex
|
||||
|
||||
t0 := time.Now()
|
||||
for _, tag := range tags {
|
||||
cmd := exec.Command("go", "run", "exp8_worker.go", tag)
|
||||
cmd.ExtraFiles = []*os.File{shmFile}
|
||||
sin, _ := cmd.StdinPipe()
|
||||
sout, _ := cmd.StdoutPipe()
|
||||
cmd.Stderr = os.Stderr
|
||||
cmd.Start()
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
dec := json.NewDecoder(bufio.NewReader(sout))
|
||||
w := bufio.NewWriter(sin)
|
||||
enc := json.NewEncoder(w)
|
||||
held := false
|
||||
for {
|
||||
var q map[string]string
|
||||
if err := dec.Decode(&q); err != nil { break }
|
||||
switch q["method"] {
|
||||
case "stage.lock": mu.Lock(); held = true
|
||||
case "stage.unlock": if held { mu.Unlock(); held = false }
|
||||
}
|
||||
cntMu.Lock(); rpcCount++; cntMu.Unlock()
|
||||
enc.Encode(map[string]bool{"ok": true}); w.Flush()
|
||||
}
|
||||
if held { mu.Unlock() }
|
||||
cmd.Wait()
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
dur := time.Since(t0)
|
||||
|
||||
off := binary.LittleEndian.Uint32(data[0:4])
|
||||
ln := binary.LittleEndian.Uint32(data[4:8])
|
||||
final := string(data[off : off+ln])
|
||||
|
||||
fmt.Printf("\n--- 结果 ---\n")
|
||||
fmt.Printf("最终 final_text 长度 = %d\n", len(final))
|
||||
counts := map[string]int{}
|
||||
for _, t := range tags { counts[t] = strings.Count(final, t) }
|
||||
fmt.Printf("各插件写入次数: %v\n", counts)
|
||||
total := 0
|
||||
for _, c := range counts { total += c }
|
||||
fmt.Printf("总字符 = %d, 长度 = %d → %s\n", total, len(final),
|
||||
map[bool]string{true:"一致 ✅ 无丢失/无撕裂", false:"不一致 ❌"}[total == len(final)])
|
||||
fmt.Printf("RPC 锁操作 = %d 次, 总耗时 %v\n", rpcCount, dur)
|
||||
fmt.Printf("\n注:写入次数少于 5×300 是 arena 64KB 上限所致(append-only 未压实),符合设计\n")
|
||||
}
|
||||
@ -0,0 +1,50 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"strconv"
|
||||
|
||||
"golang.org/x/sys/unix"
|
||||
)
|
||||
|
||||
// 模拟插件:拿锁 → 读 final_text → 追加自己的标记 → 写回 → 放锁
|
||||
// 锁通过 stdio RPC 向内核申请(方案 3.7:锁仲裁回归内核,无 cgo)
|
||||
func main() {
|
||||
tag := os.Args[1]
|
||||
shmf := os.NewFile(3, "shm")
|
||||
data, err := unix.Mmap(int(shmf.Fd()), 0, 65536, unix.PROT_READ|unix.PROT_WRITE, unix.MAP_SHARED)
|
||||
if err != nil { fmt.Fprintln(os.Stderr, "mmap:", err); os.Exit(1) }
|
||||
|
||||
dec := json.NewDecoder(bufio.NewReader(os.Stdin))
|
||||
w := bufio.NewWriter(os.Stdout)
|
||||
enc := json.NewEncoder(w)
|
||||
rpc := func(method string) {
|
||||
enc.Encode(map[string]string{"method": method}); w.Flush()
|
||||
var r map[string]interface{}; dec.Decode(&r)
|
||||
}
|
||||
|
||||
const iters = 300
|
||||
for i := 0; i < iters; i++ {
|
||||
rpc("stage.lock")
|
||||
// --- 临界区:偏移解引用读写 final_text ---
|
||||
off := binary.LittleEndian.Uint32(data[0:4])
|
||||
ln := binary.LittleEndian.Uint32(data[4:8])
|
||||
cur := string(data[off : off+ln])
|
||||
add := tag
|
||||
newS := cur + add
|
||||
// append-only arena:写到新位置
|
||||
newOff := binary.LittleEndian.Uint32(data[8:12])
|
||||
if int(newOff)+len(newS) > 65536 { rpc("stage.unlock"); break }
|
||||
copy(data[newOff:], []byte(newS))
|
||||
binary.LittleEndian.PutUint32(data[0:4], newOff)
|
||||
binary.LittleEndian.PutUint32(data[4:8], uint32(len(newS)))
|
||||
binary.LittleEndian.PutUint32(data[8:12], newOff+uint32(len(newS)))
|
||||
rpc("stage.unlock")
|
||||
}
|
||||
fmt.Fprintln(os.Stderr, "worker "+tag+" done, iters="+strconv.Itoa(iters))
|
||||
}
|
||||
60
docs/zh/experiments/plugin-arch/02-feasibility/exp9.go
Normal file
60
docs/zh/experiments/plugin-arch/02-feasibility/exp9.go
Normal file
@ -0,0 +1,60 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os/exec"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
func run(name, arg string, mu *sync.Mutex, crashed *bool) {
|
||||
cmd := exec.Command("go", "run", "exp9_worker.go", arg)
|
||||
sin, _ := cmd.StdinPipe(); sout, _ := cmd.StdoutPipe()
|
||||
cmd.Stderr = nil
|
||||
cmd.Start()
|
||||
dec := json.NewDecoder(bufio.NewReader(sout))
|
||||
w := bufio.NewWriter(sin); enc := json.NewEncoder(w)
|
||||
held := false
|
||||
for {
|
||||
var q map[string]string
|
||||
if err := dec.Decode(&q); err != nil { break }
|
||||
switch q["method"] {
|
||||
case "stage.lock": mu.Lock(); held = true; fmt.Printf(" [%s] 获得锁\n", name)
|
||||
case "stage.unlock": if held { mu.Unlock(); held = false; fmt.Printf(" [%s] 释放锁\n", name) }
|
||||
}
|
||||
enc.Encode(map[string]bool{"ok":true}); w.Flush()
|
||||
}
|
||||
err := cmd.Wait()
|
||||
// 关键:进程死了,内核侧检测到 EOF/退出 → 强制释放它持有的锁
|
||||
if held {
|
||||
mu.Unlock()
|
||||
*crashed = true
|
||||
fmt.Printf(" [%s] 进程死亡(%v),内核强制释放其持有的锁 ← 自愈\n", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
func main() {
|
||||
fmt.Println("=== 实验 9:持锁进程崩溃后的自愈(验证无需 robust pthread_mutex)===")
|
||||
var mu sync.Mutex
|
||||
crashed := false
|
||||
|
||||
fmt.Println("\n1) 插件 X 拿锁后 panic:")
|
||||
run("X", "crash", &mu, &crashed)
|
||||
|
||||
fmt.Println("\n2) 插件 Y 随后申请同一把锁:")
|
||||
done := make(chan bool, 1)
|
||||
go func() { run("Y", "normal", &mu, new(bool)); done <- true }()
|
||||
select {
|
||||
case <-done:
|
||||
fmt.Println("\n✅ Y 正常获得并释放锁 —— 无死锁")
|
||||
fmt.Println(" → 内核持有锁的所有权,进程死亡由 Wait()/EOF 检测并强制释放")
|
||||
fmt.Println(" → 不需要 PTHREAD_PROCESS_SHARED|ROBUST,也不需要处理 EOWNERDEAD")
|
||||
fmt.Println(" → 整个架构可做到零 cgo")
|
||||
case <-time.After(15 * time.Second):
|
||||
fmt.Println("\n❌ 死锁:Y 拿不到锁(说明需要 robust 语义)")
|
||||
}
|
||||
_ = crashed
|
||||
}
|
||||
@ -0,0 +1,12 @@
|
||||
//go:build ignore
|
||||
package main
|
||||
|
||||
import ("bufio";"encoding/json";"os")
|
||||
func main() {
|
||||
dec := json.NewDecoder(bufio.NewReader(os.Stdin))
|
||||
w := bufio.NewWriter(os.Stdout); enc := json.NewEncoder(w)
|
||||
rpc := func(m string) { enc.Encode(map[string]string{"method":m}); w.Flush(); var r map[string]interface{}; dec.Decode(&r) }
|
||||
rpc("stage.lock")
|
||||
if os.Args[1] == "crash" { panic("持锁时崩溃") } // 拿着锁死掉
|
||||
rpc("stage.unlock")
|
||||
}
|
||||
Reference in New Issue
Block a user