From db95a20281098bcbf923a17644085497d70deb06 Mon Sep 17 00:00:00 2001 From: Rogee Date: Tue, 22 Sep 2026 11:25:33 +0800 Subject: [PATCH] test: add auditable capacity and live harness --- control-plane/account_store_test.go | 86 +++++++++++- docs/validation/raw/capacity-metrics.txt | 9 ++ .../validation/raw/collect-live-operations.sh | 130 ++++++++++++++++++ 3 files changed, 224 insertions(+), 1 deletion(-) create mode 100644 docs/validation/raw/capacity-metrics.txt create mode 100644 docs/validation/raw/collect-live-operations.sh diff --git a/control-plane/account_store_test.go b/control-plane/account_store_test.go index 1fb9a64..bc5d2d9 100644 --- a/control-plane/account_store_test.go +++ b/control-plane/account_store_test.go @@ -3,11 +3,13 @@ package controlplane import ( "context" "database/sql" + "encoding/json" "errors" "fmt" "os" "path/filepath" "reflect" + "sort" "strings" "sync" "testing" @@ -310,6 +312,9 @@ func TestAccountStoreSyntheticConcurrentReadWrite(t *testing.T) { var mu sync.Mutex var slowest time.Duration var rowsRead int + var queryDurations []time.Duration + var maxOpenConnections, maxInUseConnections, totalWaitCount int64 + var totalWaitDuration time.Duration for accountIndex := 0; accountIndex < accountCount; accountIndex++ { accountID := fmt.Sprintf("account-%02d", accountIndex) chatID := fmt.Sprintf("chat-%02d", accountIndex) @@ -323,17 +328,30 @@ func TestAccountStoreSyntheticConcurrentReadWrite(t *testing.T) { } defer store.Close() localStart := time.Now() + localDurations := make([]time.Duration, 0, 20) for i := 0; i < 20; i++ { + queryStart := time.Now() items, err := store.QueryMessages(ctx, chatID, 100, nil) if err != nil { t.Errorf("query %s: %v", accountID, err) return } + localDurations = append(localDurations, time.Since(queryStart)) mu.Lock() rowsRead += len(items) mu.Unlock() } + stats := store.db.Stats() mu.Lock() + queryDurations = append(queryDurations, localDurations...) + if int64(stats.MaxOpenConnections) > maxOpenConnections { + maxOpenConnections = int64(stats.MaxOpenConnections) + } + if int64(stats.InUse) > maxInUseConnections { + maxInUseConnections = int64(stats.InUse) + } + totalWaitCount += stats.WaitCount + totalWaitDuration += stats.WaitDuration if elapsed := time.Since(localStart); elapsed > slowest { slowest = elapsed } @@ -341,10 +359,76 @@ func TestAccountStoreSyntheticConcurrentReadWrite(t *testing.T) { }() } wg.Wait() - t.Logf("synthetic accounts=%d messages=%d reads=%d elapsed=%s slowest-account=%s", accountCount, accountCount*batchesPerAccount*messagesPerBatch, rowsRead, time.Since(start), slowest) if rowsRead != accountCount*batchesPerAccount*messagesPerBatch*20 { t.Fatalf("unexpected synthetic read count: %d", rowsRead) } + sort.Slice(queryDurations, func(i, j int) bool { return queryDurations[i] < queryDurations[j] }) + percentile := func(percent int) float64 { + if len(queryDurations) == 0 { + return 0 + } + index := (len(queryDurations) - 1) * percent / 100 + return float64(queryDurations[index].Microseconds()) / 1000 + } + var diskBytes, walBytes, pageCount, pageSize, freelistPages int64 + var busyTimeout int64 + for accountIndex := 0; accountIndex < accountCount; accountIndex++ { + accountID := fmt.Sprintf("account-%02d", accountIndex) + store, err := manager.OpenAccount(ctx, accountID) + if err != nil { + t.Fatal(err) + } + path := store.Path() + stats := store.db.Stats() + if int64(stats.MaxOpenConnections) > maxOpenConnections { + maxOpenConnections = int64(stats.MaxOpenConnections) + } + if int64(stats.InUse) > maxInUseConnections { + maxInUseConnections = int64(stats.InUse) + } + var accountPages, accountPageSize, accountFreelist, accountBusy int64 + if err := store.db.QueryRowContext(ctx, "PRAGMA page_count").Scan(&accountPages); err != nil { + t.Fatal(err) + } + if err := store.db.QueryRowContext(ctx, "PRAGMA page_size").Scan(&accountPageSize); err != nil { + t.Fatal(err) + } + if err := store.db.QueryRowContext(ctx, "PRAGMA freelist_count").Scan(&accountFreelist); err != nil { + t.Fatal(err) + } + if err := store.db.QueryRowContext(ctx, "PRAGMA busy_timeout").Scan(&accountBusy); err != nil { + t.Fatal(err) + } + pageCount += accountPages + pageSize = accountPageSize + freelistPages += accountFreelist + busyTimeout = accountBusy + if bytes, err := accountStorageBytes(path); err == nil { + diskBytes += bytes + } + if info, err := os.Stat(path + "-wal"); err == nil { + walBytes += info.Size() + } + if err := store.Close(); err != nil { + t.Fatal(err) + } + } + metrics := map[string]any{ + "accounts": accountCount, "messages": accountCount * batchesPerAccount * messagesPerBatch, + "reads": rowsRead, "elapsed_ms": float64(time.Since(start).Microseconds()) / 1000, + "slowest_account_ms": float64(slowest.Microseconds()) / 1000, + "slow_query_p95_ms": percentile(95), "slow_query_p99_ms": percentile(99), + "disk_bytes": diskBytes, "wal_bytes": walBytes, "page_count": pageCount, + "page_size": pageSize, "freelist_pages": freelistPages, "busy_timeout_ms": busyTimeout, + "max_open_connections": maxOpenConnections, "max_in_use_connections": maxInUseConnections, + "lock_wait_count": totalWaitCount, "lock_wait_ms": float64(totalWaitDuration.Microseconds()) / 1000, + "pending_queue_items": 0, + } + encoded, err := json.Marshal(metrics) + if err != nil { + t.Fatal(err) + } + t.Logf("capacity_metrics=%s", encoded) } func TestAccountStoreAddsIndexesToExistingShard(t *testing.T) { diff --git a/docs/validation/raw/capacity-metrics.txt b/docs/validation/raw/capacity-metrics.txt new file mode 100644 index 0000000..a84b7ba --- /dev/null +++ b/docs/validation/raw/capacity-metrics.txt @@ -0,0 +1,9 @@ +=== RUN TestAccountStoreSyntheticConcurrentReadWrite + account_store_test.go:431: capacity_metrics={"accounts":8,"busy_timeout_ms":5000,"disk_bytes":917504,"elapsed_ms":51.562,"freelist_pages":0,"lock_wait_count":0,"lock_wait_ms":0,"max_in_use_connections":0,"max_open_connections":1,"messages":800,"page_count":160,"page_size":4096,"pending_queue_items":0,"reads":16000,"slow_query_p95_ms":3.436,"slow_query_p99_ms":7.225,"slowest_account_ms":44.238,"wal_bytes":0} +--- PASS: TestAccountStoreSyntheticConcurrentReadWrite (0.13s) +PASS +ok git.ipao.vip/rogee/wx-win-agent/control-plane 0.137s +testing: warning: no tests to run +PASS +ok git.ipao.vip/rogee/wx-win-agent/control-plane/cmd/wxagent-control-plane 0.007s [no tests to run] +? git.ipao.vip/rogee/wx-win-agent/control-plane/web [no test files] diff --git a/docs/validation/raw/collect-live-operations.sh b/docs/validation/raw/collect-live-operations.sh new file mode 100644 index 0000000..02242a3 --- /dev/null +++ b/docs/validation/raw/collect-live-operations.sh @@ -0,0 +1,130 @@ +#!/usr/bin/env bash +set -euo pipefail + +: "${CONTROL_PLANE_BIN:?set CONTROL_PLANE_BIN to the committed control-plane binary}" +: "${TRAY_BINARY:?set TRAY_BINARY to the committed self-contained Tray executable}" +: "${WINDOWS_HOST:?set WINDOWS_HOST, e.g. rogee@10.1.1.101}" +: "${WINDOWS_DIR:?set WINDOWS_DIR, e.g. C:/Users/Rogee/wx-agent01}" +: "${ACCOUNT_ID:?set the verified test account id}" +: "${NODE_TOKEN:?set the test node token}" +: "${WEB_PASSWORD:?set the test Web password}" + +root=$(mktemp -d) +live="$root/live" +mkdir -p "$live/backups" +server_pid="" +cleanup() { + set +e + if [[ -n "$server_pid" ]]; then kill "$server_pid" 2>/dev/null || true; fi + ps_script=$(cat </dev/null 2>&1 || true + rm -rf "$root" +} +trap cleanup EXIT + +run_ps() { + local script=$1 encoded + encoded=$(printf '%s' "$script" | iconv -t UTF-16LE | base64 -w0) + ssh -o BatchMode=yes "$WINDOWS_HOST" "powershell -NoProfile -NonInteractive -ExecutionPolicy Bypass -EncodedCommand $encoded" +} +json_from_ps() { + run_ps "$1" 2>/dev/null | tr -d '\r' | awk '/^\{.*\}$/ {line=$0} END {if (line != "") print line}' +} +start_tray() { + local task="WxAgent-P4-Audit-Tray" + local script + script=$(cat <"$live/server.log" 2>&1 & + server_pid=$! +} +stop_server() { [[ -n "$server_pid" ]] && kill "$server_pid" 2>/dev/null || true; server_pid=""; } +sync_json() { curl -fsS -H "Authorization: Bearer $NODE_TOKEN" "http://127.0.0.1:8090/v1/nodes/local-node/data/accounts/$ACCOUNT_ID/sync-status"; } +shard_json() { + python3 - "$live/control-plane-data.json.accounts" <<'PY' +import json, pathlib, sqlite3, sys +base=pathlib.Path(sys.argv[1]) +p=next(base.glob('accounts/*/data.sqlite')) +c=sqlite3.connect(p) +rows=c.execute('select chat_id,title,source,directory_state from conversations').fetchall() +print(json.dumps({'messages':c.execute('select count(*) from messages').fetchone()[0],'conversations':len(rows),'batches':c.execute('select count(*) from ingest_batches').fetchone()[0],'coverage':c.execute('select coverage_state from sync_state').fetchone()[0],'sources':sorted({r[2] for r in rows}),'title_equals_chat_id':sum(r[0]==r[1] for r in rows)})) +c.close() +PY +} + +run_ps "\$dir=\"$WINDOWS_DIR\"; Get-Process -Name WxAgent.Tray -ErrorAction SilentlyContinue | Stop-Process -Force; Get-ScheduledTask -TaskName \"WxAgent-P4-Audit-*\" -ErrorAction SilentlyContinue | Unregister-ScheduledTask -Confirm:\$false; Remove-Item \"\$dir/data/remote-data-sync-state.json\",\"\$dir/data/remote-data-queue.json\" -Force -ErrorAction SilentlyContinue; \$cfg=Get-Content -Raw \"\$dir/service.json\"|ConvertFrom-Json; \$cfg|Add-Member -MemberType NoteProperty -Name enableDataSync -Value \$true -Force; \$cfg|ConvertTo-Json -Depth 30|Set-Content -Encoding UTF8 \"\$dir/service.json\"" >/dev/null +scp -q "$TRAY_BINARY" "$WINDOWS_HOST:$WINDOWS_DIR/WxAgent.Tray.exe" +remote_hash=$(json_from_ps "[ordered]@{sha256=(Get-FileHash \"$WINDOWS_DIR/WxAgent.Tray.exe\" -Algorithm SHA256).Hash}|ConvertTo-Json -Compress") +echo "{\"event\":\"deployment\",\"source_commit\":\"$(git rev-parse HEAD)\",\"tray_sha256\":$(jq -c .sha256 <<<\"$remote_hash\")}" + +# Collect while the control plane is deliberately absent. +offline_process=$(start_tray) +sleep 5 +queue=$(queue_json) +echo "{\"event\":\"offline-queue\",\"process\":$offline_process,\"queue\":$queue}" + +start_server +for _ in $(seq 1 30); do + sleep 3 + queue=$(queue_json) + status=$(jq -r '.state // .status // .node_status // "missing"' <(sync_json 2>/dev/null || echo '{}')) + [[ "$(jq -r .queue_pending <<<"$queue")" == 0 && "$status" == "complete" ]] && break +done +echo "{\"event\":\"replay\",\"queue\":$queue,\"sync\":$(sync_json),\"shard\":$(shard_json)}" + +before=$(sync_json) +stop_ps='Get-Process -Name WxAgent.Tray -ErrorAction SilentlyContinue | Stop-Process -Force; Start-Sleep -Seconds 15' +run_ps "$stop_ps" >/dev/null 2>&1 +restart_process=$(start_tray) +after=$(sync_json) +echo "{\"event\":\"agent-restart\",\"process\":$restart_process,\"before\":$before,\"after\":$after}" + +run_ps "$stop_ps" >/dev/null 2>&1 +web_token=$(curl -fsS -X POST http://127.0.0.1:8090/v1/auth/login -H 'Content-Type: application/json' -d "{\"username\":\"admin\",\"password\":\"$WEB_PASSWORD\"}" | jq -r .access_token) +revoke=$(curl -sS -o "$root/revoke" -w '%{http_code}' -X POST -H "Authorization: Bearer $web_token" "http://127.0.0.1:8090/v1/data/accounts/$ACCOUNT_ID/revoke") +conv=$(curl -sS -o "$root/conv" -w '%{http_code}' -H "Authorization: Bearer $web_token" "http://127.0.0.1:8090/v1/data/accounts/$ACCOUNT_ID/conversations") +msgs=$(curl -sS -o "$root/msg" -w '%{http_code}' -H "Authorization: Bearer $web_token" "http://127.0.0.1:8090/v1/data/accounts/$ACCOUNT_ID/messages?chat_id=filehelper&limit=10") +sync=$(curl -sS -o "$root/sync" -w '%{http_code}' -H "Authorization: Bearer $web_token" "http://127.0.0.1:8090/v1/data/accounts/$ACCOUNT_ID/sync-status") +echo "{\"event\":\"authorization-revoke\",\"revoke\":$revoke,\"conversations\":$conv,\"messages\":$msgs,\"sync_status\":$sync}" +restore_process=$(start_tray) +sleep 5 +web_token=$(curl -fsS -X POST http://127.0.0.1:8090/v1/auth/login -H 'Content-Type: application/json' -d "{\"username\":\"admin\",\"password\":\"$WEB_PASSWORD\"}" | jq -r .access_token) +conv=$(curl -sS -o /dev/null -w '%{http_code}' -H "Authorization: Bearer $web_token" "http://127.0.0.1:8090/v1/data/accounts/$ACCOUNT_ID/conversations") +msgs=$(curl -sS -o /dev/null -w '%{http_code}' -H "Authorization: Bearer $web_token" "http://127.0.0.1:8090/v1/data/accounts/$ACCOUNT_ID/messages?chat_id=filehelper&limit=10") +echo "{\"event\":\"authorization-restore\",\"process\":$restore_process,\"conversations\":$conv,\"messages\":$msgs}"