test: add auditable capacity and live harness

This commit is contained in:
2026-09-22 11:25:33 +08:00
parent 9a304ab5e2
commit db95a20281
3 changed files with 224 additions and 1 deletions
+85 -1
View File
@@ -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) {
+9
View File
@@ -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]
@@ -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 <<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
\$cfg = Get-Content -Raw "\$dir/service.json" | ConvertFrom-Json
\$cfg | Add-Member -MemberType NoteProperty -Name enableDataSync -Value \$false -Force
\$cfg | ConvertTo-Json -Depth 30 | Set-Content -Encoding UTF8 "\$dir/service.json"
Remove-Item "\$dir/data/remote-data-sync-state.json","\$dir/data/remote-data-queue.json" -Force -ErrorAction SilentlyContinue
PS
)
run_ps "$ps_script" >/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 <<PS
\$dir = "$WINDOWS_DIR"
\$task = "$task"
Unregister-ScheduledTask -TaskName \$task -Confirm:\$false -ErrorAction SilentlyContinue | Out-Null
\$action = New-ScheduledTaskAction -Execute "\$dir/WxAgent.Tray.exe" -WorkingDirectory \$dir
\$trigger = New-ScheduledTaskTrigger -Once -At (Get-Date).AddMinutes(1)
\$principal = New-ScheduledTaskPrincipal -UserId "DESKTOP-EGI7QCK\\Rogee" -LogonType Interactive -RunLevel Highest
Register-ScheduledTask -TaskName \$task -Action \$action -Trigger \$trigger -Principal \$principal -Force | Out-Null
Start-ScheduledTask -TaskName \$task
Start-Sleep -Seconds 35
\$p = Get-Process -Name WxAgent.Tray | Select-Object -First 1
[ordered]@{pid=\$p.Id;session=\$p.SessionId}|ConvertTo-Json -Compress
PS
)
json_from_ps "$script"
}
queue_json() {
json_from_ps "\$q=\"$WINDOWS_DIR/data/remote-data-queue.json\"; \$pending=0; \$bytes=0; if(Test-Path \$q){\$bytes=(Get-Item \$q).Length; try{\$pending=@((Get-Content -Raw \$q|ConvertFrom-Json).items).Count}catch{}}; [ordered]@{queue_bytes=\$bytes;queue_pending=\$pending}|ConvertTo-Json -Compress"
}
start_server() {
WXAGENT_CONTROL_PLANE_ADDR=0.0.0.0:8090 \
WXAGENT_CONTROL_PLANE_DATA="$live/control-plane-data.json" \
WXAGENT_CONTROL_PLANE_BACKUP_DIR="$live/backups" \
WXAGENT_CONTROL_PLANE_BACKUP_INTERVAL=10s \
WXAGENT_CONTROL_PLANE_BACKUP_COUNT=3 \
WXAGENT_NODE_ID=local-node \
WXAGENT_NODE_TOKEN="$NODE_TOKEN" \
WXAGENT_WEB_USER=admin \
WXAGENT_WEB_PASSWORD="$WEB_PASSWORD" \
"$CONTROL_PLANE_BIN" >"$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}"