Initial commit: forge-tools-proxmox — MCP-сервер для Proxmox VE
This commit is contained in:
@@ -0,0 +1,354 @@
|
||||
package pve
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
)
|
||||
|
||||
// api.go — типизированные обёртки над Proxmox VE REST (api2/json).
|
||||
// Чтения возвращают raw-тело "data"; мутации — UPID (асинхронная задача),
|
||||
// чтобы хендлер вернул идентификатор, а завершение опрашивали через
|
||||
// task_status(wait) — так мы не блокируем единственный loop (§9.6).
|
||||
//
|
||||
// Путь идёт через joinURL (PathEscape), а значения идентификаторов перед
|
||||
// этим ещё валидируются guard'ом в tools-слое — двойная защита.
|
||||
|
||||
// Version — версия PVE в кластере.
|
||||
func (t *Tenant) Version(ctx context.Context, alias string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, "/version")
|
||||
}
|
||||
|
||||
// ClusterStatus — состояние кластера (quorum, узлы, версии).
|
||||
func (t *Tenant) ClusterStatus(ctx context.Context, alias string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, "/cluster/status")
|
||||
}
|
||||
|
||||
// ClusterResources — ресурсы кластера по PVE-типу (vm/storage/node/sdn).
|
||||
func (t *Tenant) ClusterResources(ctx context.Context, alias, rtype string) (json.RawMessage, error) {
|
||||
path := "/cluster/resources"
|
||||
if rtype != "" {
|
||||
path += "?type=" + url.QueryEscape(rtype)
|
||||
}
|
||||
return t.GET(ctx, alias, path)
|
||||
}
|
||||
|
||||
// GuestResources — VM/CT по типу (qemu|lxc): тянет type=vm и фильтрует по
|
||||
// полю "type" каждой записи (PVE не принимает type=qemu напрямую).
|
||||
func (t *Tenant) GuestResources(ctx context.Context, alias, kind string) (json.RawMessage, error) {
|
||||
raw, err := t.GET(ctx, alias, "/cluster/resources?type=vm")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var arr []map[string]any
|
||||
if err := json.Unmarshal(raw, &arr); err != nil {
|
||||
return raw, nil
|
||||
}
|
||||
out := make([]map[string]any, 0, len(arr))
|
||||
for _, it := range arr {
|
||||
if it["type"] == kind {
|
||||
out = append(out, it)
|
||||
}
|
||||
}
|
||||
return json.Marshal(out)
|
||||
}
|
||||
|
||||
// Nodes — список нод кластера.
|
||||
func (t *Tenant) Nodes(ctx context.Context, alias string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, "/nodes")
|
||||
}
|
||||
|
||||
// NodeStatus — детальное состояние ноды.
|
||||
func (t *Tenant) NodeStatus(ctx context.Context, alias, node string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/status"))
|
||||
}
|
||||
|
||||
// NodeNetwork — сетевые интерфейсы ноды (мосты/eth).
|
||||
func (t *Tenant) NodeNetwork(ctx context.Context, alias, node string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/network"))
|
||||
}
|
||||
|
||||
// VMConfig — конфиг QEMU-ВМ.
|
||||
func (t *Tenant) VMConfig(ctx context.Context, alias, node string, vmid int) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid), "/config"))
|
||||
}
|
||||
|
||||
// VMStatus — текущий статус QEMU-ВМ.
|
||||
func (t *Tenant) VMStatus(ctx context.Context, alias, node string, vmid int) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid), "/status/current"))
|
||||
}
|
||||
|
||||
// ContainerConfig — конфиг LXC.
|
||||
func (t *Tenant) ContainerConfig(ctx context.Context, alias, node string, vmid int) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/lxc", itoa(vmid), "/config"))
|
||||
}
|
||||
|
||||
// ContainerStatus — текущий статус LXC.
|
||||
func (t *Tenant) ContainerStatus(ctx context.Context, alias, node string, vmid int) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/lxc", itoa(vmid), "/status/current"))
|
||||
}
|
||||
|
||||
// StorageList — хранилища кластера.
|
||||
func (t *Tenant) StorageList(ctx context.Context, alias string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, "/storage")
|
||||
}
|
||||
|
||||
// NodeStorage — хранилища конкретной ноды.
|
||||
func (t *Tenant) NodeStorage(ctx context.Context, alias, node string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/storage"))
|
||||
}
|
||||
|
||||
// StorageContent — содержимое хранилища (iso/vztmpl/backup/images).
|
||||
func (t *Tenant) StorageContent(ctx context.Context, alias, node, storage, contentType string) (json.RawMessage, error) {
|
||||
path := joinURL("/nodes", node, "/storage", storage, "/content")
|
||||
if contentType != "" {
|
||||
path += "?content=" + url.QueryEscape(contentType)
|
||||
}
|
||||
return t.GET(ctx, alias, path)
|
||||
}
|
||||
|
||||
// SnapshotList — снапшоты VM/CT.
|
||||
func (t *Tenant) SnapshotList(ctx context.Context, alias, node, vmtype string, vmid int) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, vmtype, itoa(vmid), "/snapshot"))
|
||||
}
|
||||
|
||||
// BackupList — задачи vzdump (по ноде/всему кластеру).
|
||||
func (t *Tenant) BackupList(ctx context.Context, alias, node string) (json.RawMessage, error) {
|
||||
if node != "" {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/tasks")+"?typefilter=vzdump")
|
||||
}
|
||||
return t.GET(ctx, alias, "/cluster/tasks?typefilter=vzdump")
|
||||
}
|
||||
|
||||
// TasksList — последние задачи кластера.
|
||||
func (t *Tenant) TasksList(ctx context.Context, alias, node string, limit int) (json.RawMessage, error) {
|
||||
if node != "" {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/tasks")+"?limit="+itoa(limit))
|
||||
}
|
||||
return t.GET(ctx, alias, "/cluster/tasks?limit="+itoa(limit))
|
||||
}
|
||||
|
||||
// TaskStatus — статус задачи по UPID.
|
||||
func (t *Tenant) TaskStatus(ctx context.Context, alias, node, upid string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/tasks", upid, "/status"))
|
||||
}
|
||||
|
||||
// TaskLog — лог задачи по UPID.
|
||||
func (t *Tenant) TaskLog(ctx context.Context, alias, node, upid string, limit int) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/tasks", upid, "/log")+"?limit="+itoa(limit))
|
||||
}
|
||||
|
||||
// GuestIPs — IP-адреса гостя через QEMU guest-agent (только QEMU, agent:1).
|
||||
func (t *Tenant) GuestIPs(ctx context.Context, alias, node string, vmid int) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid), "/agent/network-get-interfaces"))
|
||||
}
|
||||
|
||||
// --- Мутации → возвращают UPID или пустую строку (у части POST нет upid) ---
|
||||
|
||||
// VMStart / VMStop / VMReboot / VMShutdown / VMSuspend / VMResume — lifecycle.
|
||||
// graceful может быть 0 (обычный), 1/2/... для shutdown. Здесь передаём action.
|
||||
func (t *Tenant) VMAction(ctx context.Context, alias, node string, vmid int, action string, form url.Values) (string, error) {
|
||||
if form == nil {
|
||||
form = url.Values{}
|
||||
}
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid), "/status", action), form)
|
||||
}
|
||||
|
||||
// ContainerAction — lifecycle LXC (start/stop/reboot/shutdown).
|
||||
func (t *Tenant) ContainerAction(ctx context.Context, alias, node string, vmid int, action string) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, "/lxc", itoa(vmid), "/status", action), nil)
|
||||
}
|
||||
|
||||
// VMClone — клонирование из template/vmid.
|
||||
func (t *Tenant) VMClone(ctx context.Context, alias, node string, vmid, newid int, form url.Values) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid), "/clone"), form)
|
||||
}
|
||||
|
||||
// ContainerClone — клонирование LXC.
|
||||
func (t *Tenant) ContainerClone(ctx context.Context, alias, node string, vmid, newid int, form url.Values) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, "/lxc", itoa(vmid), "/clone"), form)
|
||||
}
|
||||
|
||||
// VMDelete — удаление QEMU-ВМ.
|
||||
func (t *Tenant) VMDelete(ctx context.Context, alias, node string, vmid int) (string, error) {
|
||||
return t.DELETEUPID(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid)))
|
||||
}
|
||||
|
||||
// ContainerDelete — удаление LXC.
|
||||
func (t *Tenant) ContainerDelete(ctx context.Context, alias, node string, vmid int) (string, error) {
|
||||
return t.DELETEUPID(ctx, alias, joinURL("/nodes", node, "/lxc", itoa(vmid)))
|
||||
}
|
||||
|
||||
// VMConvertTemplate — превращает ВМ в шаблон.
|
||||
func (t *Tenant) VMConvertTemplate(ctx context.Context, alias, node string, vmid int) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid), "/template"), nil)
|
||||
}
|
||||
|
||||
// NextID — следующий свободный VMID в кластере.
|
||||
func (t *Tenant) NextID(ctx context.Context, alias string) (json.RawMessage, error) {
|
||||
return t.GET(ctx, alias, "/cluster/nextid")
|
||||
}
|
||||
|
||||
// SnapshotCreate / SnapshotDelete / SnapshotRollback — снапшоты VM/CT.
|
||||
func (t *Tenant) SnapshotCreate(ctx context.Context, alias, node, vmtype string, vmid int, form url.Values) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, vmtype, itoa(vmid), "/snapshot"), form)
|
||||
}
|
||||
|
||||
func (t *Tenant) SnapshotDelete(ctx context.Context, alias, node, vmtype string, vmid int, name string) (string, error) {
|
||||
return t.DELETEUPID(ctx, alias, joinURL("/nodes", node, vmtype, itoa(vmid), "/snapshot", name))
|
||||
}
|
||||
|
||||
func (t *Tenant) SnapshotRollback(ctx context.Context, alias, node, vmtype string, vmid int, name string) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, vmtype, itoa(vmid), "/snapshot", name, "/rollback"), nil)
|
||||
}
|
||||
|
||||
// BackupCreate — запускает vzdump (async UPID).
|
||||
func (t *Tenant) BackupCreate(ctx context.Context, alias, node, vmtype string, vmid int, form url.Values) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, "/vzdump"), withBackupForm(form, vmtype, vmid))
|
||||
}
|
||||
|
||||
// BackupRestore (vzdump restore) реализуется через clone из backup volume;
|
||||
// в v1 это отдельный инструмент, см. tools/backup.go.
|
||||
|
||||
// TaskStop — отменяет/стирает задачу.
|
||||
func (t *Tenant) TaskStop(ctx context.Context, alias, node, upid string) (string, error) {
|
||||
return t.DELETEUPID(ctx, alias, joinURL("/nodes", node, "/tasks", upid))
|
||||
}
|
||||
|
||||
// ConfigPost — POST на /config (VM/CT): изменение конфига, включая
|
||||
// добавление/удаление дисков и сетевых интерфейсов (структурный путь).
|
||||
func (t *Tenant) ConfigPost(ctx context.Context, alias, node, vmtype string, vmid int, form url.Values) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, vmtype, itoa(vmid), "/config"), form)
|
||||
}
|
||||
|
||||
// ResizeDisk — изменение размера диска VM/CT (/resize).
|
||||
func (t *Tenant) ResizeDisk(ctx context.Context, alias, node, vmtype string, vmid int, form url.Values) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, vmtype, itoa(vmid), "/resize"), form)
|
||||
}
|
||||
|
||||
// MoveDisk — перенос диска (storage-миграция) на другую подсистему хранения.
|
||||
func (t *Tenant) MoveDisk(ctx context.Context, alias, node string, vmid int, form url.Values) (string, error) {
|
||||
return t.POSTUPID(ctx, alias, joinURL("/nodes", node, "/qemu", itoa(vmid), "/move_disk"), form)
|
||||
}
|
||||
|
||||
// TaskWait ожидает завершения задачи по UPID (с капом task_poll_max_sec),
|
||||
// возвращая финальный статус. Poll строго bounded (§9.6) — не блокируем
|
||||
// loop дольше капа даже для долгой операции.
|
||||
func (t *Tenant) TaskWait(ctx context.Context, alias, node, upid string) (string, error) {
|
||||
capSec := t.cfg.TaskPollMaxSec
|
||||
if capSec <= 0 {
|
||||
capSec = DefaultTaskPollMaxSec
|
||||
}
|
||||
deadline := time.Now().Add(time.Duration(capSec) * time.Second)
|
||||
for {
|
||||
status, err := t.TaskStatus(ctx, alias, node, upid)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
done := taskIsDone(status)
|
||||
if done {
|
||||
return pretty(json.RawMessage(status)), nil
|
||||
}
|
||||
if time.Now().After(deadline) {
|
||||
return "", fmt.Errorf("task %s still running after %ds (give timeout)", upid, capSec)
|
||||
}
|
||||
if !sleep(ctx, 2*time.Second) {
|
||||
return "", fmt.Errorf("task wait interrupted: %w", ctx.Err())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// taskIsDone сообщает, завершилась ли задача (по статусам PVE).
|
||||
func taskIsDone(raw json.RawMessage) bool {
|
||||
var obj struct {
|
||||
Status string `json:"status"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &obj); err != nil {
|
||||
return true
|
||||
}
|
||||
switch obj.Status {
|
||||
case "stopped", "failed", "error": // completed (либо упала)
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// --- низкоуровневый доступ (GET/POST/DELETE) ---
|
||||
|
||||
// GET — чтение "data".
|
||||
func (t *Tenant) GET(ctx context.Context, alias, path string) (json.RawMessage, error) {
|
||||
c, err := t.Client(ctx, alias)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return c.Get(ctx, path)
|
||||
}
|
||||
|
||||
// POSTUPID — POST и возврат UPID из ответа ("upid" либо null).
|
||||
func (t *Tenant) POSTUPID(ctx context.Context, alias, path string, form url.Values) (string, error) {
|
||||
c, err := t.Client(ctx, alias)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
data, err := c.Post(ctx, path, form)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return upidOf(data), nil
|
||||
}
|
||||
|
||||
// DELETEUPID — DELETE и возврат UPID из ответа.
|
||||
func (t *Tenant) DELETEUPID(ctx context.Context, alias, path string) (string, error) {
|
||||
c, err := t.Client(ctx, alias)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
data, err := c.Delete(ctx, path)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return upidOf(data), nil
|
||||
}
|
||||
|
||||
// upidOf вытаскивает "upid" из JSON-объекта данных (у части POST его нет).
|
||||
func upidOf(data json.RawMessage) string {
|
||||
var obj struct {
|
||||
UPID string `json:"upid"`
|
||||
}
|
||||
if err := json.Unmarshal(data, &obj); err == nil && obj.UPID != "" {
|
||||
return obj.UPID
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// withBackupForm добавляет типы/цель в form vzdump.
|
||||
func withBackupForm(form url.Values, vmtype string, vmid int) url.Values {
|
||||
if form == nil {
|
||||
form = url.Values{}
|
||||
}
|
||||
if vmtype == "qemu" {
|
||||
form.Set("vmid", itoa(vmid))
|
||||
form.Set("mode", "snapshot")
|
||||
} else {
|
||||
form.Set("vmid", itoa(vmid))
|
||||
form.Set("mode", "suspend")
|
||||
}
|
||||
return form
|
||||
}
|
||||
|
||||
func itoa(i int) string { return fmt.Sprintf("%d", i) }
|
||||
|
||||
// pretty форматирует raw JSON для отдачи модели.
|
||||
func pretty(raw json.RawMessage) string {
|
||||
if len(raw) == 0 || string(raw) == "null" {
|
||||
return "(no data)"
|
||||
}
|
||||
var b bytes.Buffer
|
||||
if err := json.Indent(&b, raw, "", " "); err != nil {
|
||||
return string(raw)
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
@@ -0,0 +1,230 @@
|
||||
package pve
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Client — тонкий HTTP-клиент к одному Proxmox-кластеру (api2/json).
|
||||
// Только stdio-модуль (forge-tools): никакого TCP/HTTP-сервера, никакого
|
||||
// session-pool — обычный обозреватель API. Аутентификация — API-токен
|
||||
// (отзываемый, ревизуемый; не пароль-ticket с 3-сек. задержкой на 401).
|
||||
type Client struct {
|
||||
base string // полный URL, напр. https://host:8006/api2/json
|
||||
tokenID string
|
||||
secret string
|
||||
http *http.Client
|
||||
maxBody int
|
||||
}
|
||||
|
||||
// APIError — ошибка со стороны Proxmox (HTTP >= 400): доменный отказ.
|
||||
// Такие ошибки хендлеры возвращают как errorResult (модель видит и может
|
||||
// исправить), а НЕ как Go-ошибку (§3.3 контракт).
|
||||
type APIError struct {
|
||||
Status int
|
||||
Method string
|
||||
Path string
|
||||
Message string
|
||||
}
|
||||
|
||||
func (e *APIError) Error() string {
|
||||
return fmt.Sprintf("proxmox api %s %s: %s (status %d)", e.Method, e.Path, e.Message, e.Status)
|
||||
}
|
||||
|
||||
// IsAPIError сообщает, является ли ошибка доменным отказом Proxmox.
|
||||
func IsAPIError(err error) bool {
|
||||
var ae *APIError
|
||||
return errors.As(err, &ae)
|
||||
}
|
||||
|
||||
// NewClient строит клиент по host-конфигу. TLS: предпочтителен CA-файл
|
||||
// (самоподписанный кластер) — insecure только для dev-lab, не по умолчанию.
|
||||
// URL объявляет оператор, поэтому SSRF-вектора «модель ввела хост» нет.
|
||||
func NewClient(h HostConfig) (*Client, error) {
|
||||
if !isValidHTTPURL(h.URL) {
|
||||
return nil, fmt.Errorf("proxmox: invalid url for host %q", h.Alias)
|
||||
}
|
||||
tlsCfg, err := tlsConfig(h)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
transport := &http.Transport{TLSClientConfig: tlsCfg}
|
||||
return &Client{
|
||||
base: strings.TrimRight(h.URL, "/"),
|
||||
tokenID: h.TokenID,
|
||||
secret: h.TokenSecret,
|
||||
http: &http.Client{Transport: transport},
|
||||
maxBody: DefaultMaxOutputBytes,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// tlsConfig собирает конфигурацию TLS: CA-файл (рекомендован) либо
|
||||
// InsecureSkipVerify (только dev-lab). По умолчанию — системные корни
|
||||
// (fail-closed: самоподписанный сертификат не пройдёт без явного выбора).
|
||||
func tlsConfig(h HostConfig) (*tls.Config, error) {
|
||||
if h.CAFile != "" {
|
||||
pem, err := os.ReadFile(h.CAFile)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("proxmox: read ca_file: %w", err)
|
||||
}
|
||||
pool := x509.NewCertPool()
|
||||
if !pool.AppendCertsFromPEM(pem) {
|
||||
return nil, fmt.Errorf("proxmox: no certs parsed from ca_file %q", h.CAFile)
|
||||
}
|
||||
return &tls.Config{RootCAs: pool}, nil
|
||||
}
|
||||
// Оператор явно выбрал insecure — это dev-lab (самоподписанный PVE).
|
||||
return &tls.Config{InsecureSkipVerify: h.Insecure}, nil
|
||||
}
|
||||
|
||||
// Get выполняет GET и возвращает поле "data" из ответа Proxmox (raw JSON).
|
||||
func (c *Client) Get(ctx context.Context, path string) (json.RawMessage, error) {
|
||||
return c.do(ctx, http.MethodGet, path, nil)
|
||||
}
|
||||
|
||||
// Post выполняет POST с form-телом (PVE принимает application/x-www-form-urlencoded).
|
||||
func (c *Client) Post(ctx context.Context, path string, values url.Values) (json.RawMessage, error) {
|
||||
return c.do(ctx, http.MethodPost, path, values)
|
||||
}
|
||||
|
||||
// Delete выполняет DELETE.
|
||||
func (c *Client) Delete(ctx context.Context, path string) (json.RawMessage, error) {
|
||||
return c.do(ctx, http.MethodDelete, path, nil)
|
||||
}
|
||||
|
||||
// do — единая точка запроса: auth-заголовок, таймаут через ctx, retry для
|
||||
// идемпотентных GET, разбор {"data":...}, классификация APIError.
|
||||
func (c *Client) do(ctx context.Context, method, path string, form url.Values) (json.RawMessage, error) {
|
||||
// Ретраим только GET (идемпотентный) на 429/502/503/504 и сетевых сбоях.
|
||||
if method == http.MethodGet {
|
||||
var last error
|
||||
for attempt := 0; attempt < 3; attempt++ {
|
||||
data, err := c.once(ctx, method, path, form)
|
||||
if err == nil || !retryable(err) {
|
||||
return data, err
|
||||
}
|
||||
last = err
|
||||
if !sleep(ctx, backoff(attempt)) {
|
||||
return nil, last
|
||||
}
|
||||
}
|
||||
return nil, last
|
||||
}
|
||||
return c.once(ctx, method, path, form)
|
||||
}
|
||||
|
||||
// once выполняет один HTTP-запрос и разбирает ответ.
|
||||
func (c *Client) once(ctx context.Context, method, path string, form url.Values) (json.RawMessage, error) {
|
||||
var body io.Reader
|
||||
if form != nil {
|
||||
body = strings.NewReader(form.Encode())
|
||||
}
|
||||
req, err := http.NewRequestWithContext(ctx, method, c.base+path, body)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("proxmox: build request: %w", err)
|
||||
}
|
||||
req.Header.Set("Authorization", "PVEAPIToken="+c.tokenID+"="+c.secret)
|
||||
req.Header.Set("Accept", "application/json")
|
||||
if form != nil {
|
||||
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
|
||||
}
|
||||
|
||||
resp, err := c.http.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("proxmox: %s %s: %w", method, path, err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
raw, err := io.ReadAll(io.LimitReader(resp.Body, int64(c.maxBody)+1))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("proxmox: read %s %s: %w", method, path, err)
|
||||
}
|
||||
if len(raw) > c.maxBody {
|
||||
return nil, fmt.Errorf("proxmox: %s %s response exceeds %d bytes", method, path, c.maxBody)
|
||||
}
|
||||
|
||||
// Доменный отказ (>=400) — APIError с телом/текстом для модели.
|
||||
if resp.StatusCode >= 400 {
|
||||
return nil, &APIError{
|
||||
Status: resp.StatusCode,
|
||||
Method: method,
|
||||
Path: path,
|
||||
Message: apiErrorMessage(raw, resp.Status),
|
||||
}
|
||||
}
|
||||
|
||||
// PVE всегда оборачивает успех в {"data": ...}; вытаскиваем его.
|
||||
return unwrapData(raw)
|
||||
}
|
||||
|
||||
// unwrapData достаёт поле "data" из ответа {"data": ...}. Если его нет —
|
||||
// возвращаем null (напр. "undefined" у части POST).
|
||||
func unwrapData(raw []byte) (json.RawMessage, error) {
|
||||
if len(raw) == 0 {
|
||||
return json.RawMessage("null"), nil
|
||||
}
|
||||
var wrapper struct {
|
||||
Data json.RawMessage `json:"data"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &wrapper); err != nil {
|
||||
return json.RawMessage("null"), nil
|
||||
}
|
||||
if wrapper.Data == nil {
|
||||
return json.RawMessage("null"), nil
|
||||
}
|
||||
return wrapper.Data, nil
|
||||
}
|
||||
|
||||
// apiErrorMessage извлекает человекочитаемое сообщение из тела ошибки PVE.
|
||||
func apiErrorMessage(raw []byte, status string) string {
|
||||
var e struct {
|
||||
Errors map[string]string `json:"errors"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &e); err == nil && len(e.Errors) > 0 {
|
||||
var parts []string
|
||||
for k, v := range e.Errors {
|
||||
parts = append(parts, k+": "+v)
|
||||
}
|
||||
return strings.Join(parts, "; ")
|
||||
}
|
||||
s := strings.TrimSpace(string(raw))
|
||||
if s == "" || s == "null" {
|
||||
return status
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
// retryable сообщает, стоит ли повторять запрос. Доменные 429/502/503/504 —
|
||||
// да; отмену/deadline — нет (уважаем ctx).
|
||||
func retryable(err error) bool {
|
||||
var ae *APIError
|
||||
if errors.As(err, &ae) {
|
||||
return ae.Status == http.StatusTooManyRequests || ae.Status == 502 || ae.Status == 503 || ae.Status == 504
|
||||
}
|
||||
return !errors.Is(err, context.Canceled) && !errors.Is(err, context.DeadlineExceeded)
|
||||
}
|
||||
|
||||
// backoff — простой джиттер-бэкфол (0.5s, 1s).
|
||||
func backoff(attempt int) time.Duration {
|
||||
return time.Duration(500*(1<<attempt)) * time.Millisecond
|
||||
}
|
||||
|
||||
// sleep с уважением к ctx.
|
||||
func sleep(ctx context.Context, d time.Duration) bool {
|
||||
select {
|
||||
case <-time.After(d):
|
||||
return true
|
||||
case <-ctx.Done():
|
||||
return false
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,324 @@
|
||||
// Package pve — доменный слой forge-tools-proxmox: чтение per-agent
|
||||
// конфигурации (configreload / live-reload), тонкий stdio-безопасный
|
||||
// HTTP-клиент к Proxmox VE API (api2/json), guard-валидация идентификаторов
|
||||
// и потокобезопасный Manager. Здесь НЕТ MCP-зависимостей (см.
|
||||
// forge-tools/ARCHITECTURE.md §2, столп разделения домен/инструменты).
|
||||
package pve
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
"git.totmin.ru/en2zmax/forge-toolkit"
|
||||
)
|
||||
|
||||
// Политика доступа по умолчанию (least-privilege): все операции чтения
|
||||
// разрешены, мутации — только явно (read_only=false + allowlist.vmids).
|
||||
// Ниже — значения по умолчанию и безопасные границы.
|
||||
|
||||
const (
|
||||
// DefaultTimeoutSec — таймаут одного HTTP-запроса к API (§9.6 bounded).
|
||||
DefaultTimeoutSec = 15
|
||||
// DefaultTaskPollMaxSec — максимальное время ожидания завершения task
|
||||
// (UPID) в одном вызове task_status(wait=true): не блокируем loop дольше капа.
|
||||
DefaultTaskPollMaxSec = 600
|
||||
// DefaultMaxOutputBytes — кап размера тела ответа, чтобы не отдавать модели
|
||||
// гигантские JSON (лимиты больших результатов §9.6).
|
||||
DefaultMaxOutputBytes = 4 << 20 // 4 MiB
|
||||
)
|
||||
|
||||
// DefaultDenyConfigKeys — поля VM/CT config, запрещённые для изменения через
|
||||
// vm_config_update / container_config_update. Это структурно-опасные ключи:
|
||||
// их изменение надо делать выделенными инструментами + confirm, а не через
|
||||
// общий update (иначе модель может «незаметно» перестроить машину).
|
||||
var DefaultDenyConfigKeys = []string{
|
||||
"delete", "revert", "hotplug", "spice",
|
||||
"hostpci_mapping", "realm", "bootorder",
|
||||
"sockets", "cores", // меняем явно через отдельные поля, не через update
|
||||
}
|
||||
|
||||
// HostConfig — одно подключение к Proxmox-кластеру (или ноде). URL — полный
|
||||
// базовый путь API, включая /api2/json. Секреты (token_secret) берутся из
|
||||
// ${VAR} или gitignored *.local.json — никогда не коммитятся (§8).
|
||||
//
|
||||
// AllowNodes/AllowVMIDs — per-host ограничения мутаций. Их авторитетность
|
||||
// выше глобального Config.Allowlist: в мульти-гипервизорной конфигурации
|
||||
// права каждого хоста изолированы, а VMID/ноды на разных гипервизорах могут
|
||||
// пересекаться, поэтому «голый» VMID здесь НЕ идентифицирует ресурс —
|
||||
// ресурс всегда (host, node, vmid).
|
||||
type HostConfig struct {
|
||||
Alias string `json:"alias"`
|
||||
URL string `json:"url"`
|
||||
TokenID string `json:"token_id"`
|
||||
TokenSecret string `json:"token_secret"`
|
||||
CAFile string `json:"ca_file"`
|
||||
Insecure bool `json:"insecure"`
|
||||
// AllowNodes — ноды этого хоста, доступные для мутаций (fail-closed:
|
||||
// пусто = мутации уровня ноды запрещены).
|
||||
AllowNodes []string `json:"allow_nodes"`
|
||||
// AllowVMIDs — VMID/CTID этого хоста, доступные для мутаций (fail-closed:
|
||||
// пусто = мутации гостей запрещены). Ключевой перенос: права per-host.
|
||||
AllowVMIDs []int `json:"allow_vmids"`
|
||||
}
|
||||
|
||||
// Allowlist — ограничение ресурсов, доступных для МУТАЦИЙ. Пустой список =
|
||||
// fail-closed (мутации запрещены), а не «всё разрешено». Это второй слой
|
||||
// безопасности поверх least-privilege токена (§9.5).
|
||||
type Allowlist struct {
|
||||
Nodes []string `json:"nodes"`
|
||||
VMIDs []int `json:"vmids"`
|
||||
}
|
||||
|
||||
// Config — per-agent политика подключения и лимитов.
|
||||
type Config struct {
|
||||
Hosts []HostConfig `json:"hosts"`
|
||||
Default string `json:"default"`
|
||||
ReadOnly *bool `json:"read_only"`
|
||||
Allowlist Allowlist `json:"allowlist"`
|
||||
TimeoutSec int `json:"timeout_sec"`
|
||||
TaskPollMaxSec int `json:"task_poll_max_sec"`
|
||||
MaxOutputBytes int `json:"max_output_bytes"`
|
||||
denyConfigKeys []string
|
||||
}
|
||||
|
||||
// ParseConfig разбирает байты pve.json (${VAR} + валидация + дефолты).
|
||||
// Выделена отдельной функцией, чтобы её использовали и configreload.Loader
|
||||
// (live-reload), и стартовый -config. Fail-closed: без hosts — ошибка.
|
||||
func ParseConfig(data []byte) (*Config, error) {
|
||||
// ${VAR} разворачиваем до unmarshal; отсутствующая переменная → пустая
|
||||
// строка, а валидация ниже отвергнет пустой обязательный secret (не
|
||||
// подставляем мусор).
|
||||
expanded := toolkit.Expand(data)
|
||||
|
||||
cfg := &Config{}
|
||||
if err := json.Unmarshal(expanded, cfg); err != nil {
|
||||
return nil, fmt.Errorf("parse proxmox config: %w", err)
|
||||
}
|
||||
if err := cfg.normalize(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
// LoadConfig читает и парсит файл конфига по пути.
|
||||
func LoadConfig(path string) (*Config, error) {
|
||||
data, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read proxmox config %s: %w", path, err)
|
||||
}
|
||||
return ParseConfig(data)
|
||||
}
|
||||
|
||||
// normalize проверяет обязательные поля и применяет дефолты. Fail-closed:
|
||||
// невалидная политика — ошибка, а не «предположим что-то разумное».
|
||||
// Мульти-гипервизор: глобальный allowlist запрещён (права обязаны быть
|
||||
// per-host — иначе VMID пересекутся между кластерами), алиасы/URL уникальны.
|
||||
func (c *Config) normalize() error {
|
||||
if len(c.Hosts) == 0 {
|
||||
return errors.New("proxmox config: at least one host is required")
|
||||
}
|
||||
multi := len(c.Hosts) > 1
|
||||
if multi && (len(c.Allowlist.VMIDs) > 0 || len(c.Allowlist.Nodes) > 0) {
|
||||
return errors.New("proxmox config: global allowlist is not allowed with multiple hosts — set allow_vmids/allow_nodes per host")
|
||||
}
|
||||
seenAlias := map[string]bool{}
|
||||
seenURL := map[string]bool{}
|
||||
for i := range c.Hosts {
|
||||
h := &c.Hosts[i]
|
||||
if h.Alias == "" {
|
||||
return fmt.Errorf("proxmox config: host[%d].alias is required", i)
|
||||
}
|
||||
if seenAlias[h.Alias] {
|
||||
return fmt.Errorf("proxmox config: duplicate host alias %q", h.Alias)
|
||||
}
|
||||
seenAlias[h.Alias] = true
|
||||
if !isValidHTTPURL(h.URL) {
|
||||
return fmt.Errorf("proxmox config: host %q has invalid/unsupported url", h.Alias)
|
||||
}
|
||||
if seenURL[h.URL] {
|
||||
return fmt.Errorf("proxmox config: duplicate host url %q", h.URL)
|
||||
}
|
||||
seenURL[h.URL] = true
|
||||
if h.TokenID == "" || h.TokenSecret == "" {
|
||||
return fmt.Errorf("proxmox config: host %q requires token_id and token_secret", h.Alias)
|
||||
}
|
||||
}
|
||||
if c.Default == "" {
|
||||
c.Default = c.Hosts[0].Alias
|
||||
}
|
||||
if c.ReadOnly == nil {
|
||||
t := true
|
||||
c.ReadOnly = &t
|
||||
}
|
||||
if c.TimeoutSec <= 0 {
|
||||
c.TimeoutSec = DefaultTimeoutSec
|
||||
}
|
||||
if c.TaskPollMaxSec <= 0 {
|
||||
c.TaskPollMaxSec = DefaultTaskPollMaxSec
|
||||
}
|
||||
if c.MaxOutputBytes <= 0 {
|
||||
c.MaxOutputBytes = DefaultMaxOutputBytes
|
||||
}
|
||||
if len(c.denyConfigKeys) == 0 {
|
||||
c.denyConfigKeys = DefaultDenyConfigKeys
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// HasHosts сообщает, настроено ли хотя бы одно подключение.
|
||||
func (c *Config) HasHosts() bool { return c != nil && len(c.Hosts) > 0 }
|
||||
|
||||
// IsReadOnly сообщает, запрещены ли мутации (дефолт: true).
|
||||
func (c *Config) IsReadOnly() bool {
|
||||
if c == nil || c.ReadOnly == nil {
|
||||
return true
|
||||
}
|
||||
return *c.ReadOnly
|
||||
}
|
||||
|
||||
// Host возвращает подключение по алиасу (или default).
|
||||
func (c *Config) Host(alias string) (HostConfig, bool) {
|
||||
if c == nil {
|
||||
return HostConfig{}, false
|
||||
}
|
||||
if alias == "" {
|
||||
alias = c.Default
|
||||
}
|
||||
for _, h := range c.Hosts {
|
||||
if h.Alias == alias {
|
||||
return h, true
|
||||
}
|
||||
}
|
||||
return HostConfig{}, false
|
||||
}
|
||||
|
||||
// HostExists сообщает, настроен ли хост по алиасу.
|
||||
func (c *Config) HostExists(alias string) bool {
|
||||
if c == nil {
|
||||
return false
|
||||
}
|
||||
_, ok := c.Host(alias)
|
||||
return ok
|
||||
}
|
||||
|
||||
// MultiHost сообщает, настроено ли больше одного гипервизора.
|
||||
func (c *Config) MultiHost() bool { return c != nil && len(c.Hosts) > 1 }
|
||||
|
||||
// WriteAllowed решает, разрешена ли МУТАЦИЯ над VM/CT на конкретном хосте.
|
||||
// Перенос on per-host allow_vmids; глобальный allowlist — только фолбэк для
|
||||
// одно-гипервизорной конфигурации (в мульти-конфиге он запрещён). Fail-closed:
|
||||
// оба пустые → запрещено. Это исключает коллизию VMID между гипервизорами.
|
||||
func (c *Config) WriteAllowed(alias string, vmid int) bool {
|
||||
if c == nil || c.IsReadOnly() {
|
||||
return false
|
||||
}
|
||||
h, ok := c.Host(alias)
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
arr := h.AllowVMIDs
|
||||
if len(arr) == 0 {
|
||||
arr = c.Allowlist.VMIDs
|
||||
}
|
||||
return containsInt(arr, vmid)
|
||||
}
|
||||
|
||||
// NodeWriteAllowed — то же для мутаций уровня ноды, перенос on allow_nodes.
|
||||
func (c *Config) NodeWriteAllowed(alias, node string) bool {
|
||||
if c == nil || c.IsReadOnly() {
|
||||
return false
|
||||
}
|
||||
h, ok := c.Host(alias)
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
arr := h.AllowNodes
|
||||
if len(arr) == 0 {
|
||||
arr = c.Allowlist.Nodes
|
||||
}
|
||||
return containsStr(arr, node)
|
||||
}
|
||||
|
||||
func containsInt(arr []int, v int) bool {
|
||||
for _, x := range arr {
|
||||
if x == v {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func containsStr(arr []string, v string) bool {
|
||||
for _, x := range arr {
|
||||
if x == v {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// DenyConfigKey сообщает, запрещён ли ключ конфига для update-инструментов.
|
||||
func (c *Config) DenyConfigKey(key string) bool {
|
||||
if c == nil {
|
||||
return false
|
||||
}
|
||||
for _, k := range c.denyConfigKeys {
|
||||
if k == key {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// ValidateIdentifier — экспортированный guard против path-traversal для
|
||||
// идентификаторов (node, vmid, snapname, storage, upid), попадающих в URL.
|
||||
func ValidateIdentifier(s string) error { return validatePathToken(s) }
|
||||
|
||||
// ValidateUPID — guard для значений UPID (task): они содержат ':' '@' '!',
|
||||
// поэтому допустимая шире, но строго БЕЗ разделителей пути и подъёма '..'.
|
||||
func ValidateUPID(s string) error {
|
||||
if s == "" {
|
||||
return errors.New("empty task upid")
|
||||
}
|
||||
if strings.ContainsAny(s, "/\\\x00") || strings.Contains(s, "..") {
|
||||
return fmt.Errorf("unsafe upid %q", s)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// validatePathToken — guard против path-traversal: идентификаторы (node,
|
||||
// vmid, snapname, storage, upid), попадающие в URL-путь, обязаны быть из
|
||||
// безопасного алфавита и не содержать сепараторов/подъёма (анти-инъекция §9).
|
||||
func validatePathToken(s string) error {
|
||||
if s == "" {
|
||||
return errors.New("empty path identifier")
|
||||
}
|
||||
if strings.ContainsAny(s, "/\\\x00") || strings.Contains(s, "..") {
|
||||
return fmt.Errorf("unsafe path identifier %q", s)
|
||||
}
|
||||
for _, r := range s {
|
||||
switch {
|
||||
case r >= 'a' && r <= 'z':
|
||||
case r >= 'A' && r <= 'Z':
|
||||
case r >= '0' && r <= '9':
|
||||
case r == '.' || r == '_' || r == '-':
|
||||
default:
|
||||
return fmt.Errorf("unsafe path identifier %q: invalid char %q", s, r)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// isValidHTTPURL проверяет схему и наличие хоста (модель не задаёт URL —
|
||||
// его объявляет оператор; здесь лишь отсекаем явный мусор).
|
||||
func isValidHTTPURL(raw string) bool {
|
||||
if !strings.HasPrefix(raw, "http://") && !strings.HasPrefix(raw, "https://") {
|
||||
return false
|
||||
}
|
||||
rest := strings.TrimPrefix(strings.TrimPrefix(raw, "https://"), "http://")
|
||||
// хост обязан быть, но может содержать порт; путь — /api2/json или глубже.
|
||||
return rest != "" && rest != "/" && !strings.HasPrefix(rest, "/")
|
||||
}
|
||||
@@ -0,0 +1,209 @@
|
||||
package pve
|
||||
|
||||
import (
|
||||
"testing"
|
||||
)
|
||||
|
||||
// include: unit-тесты политики/конфига домена (без сети и без MCP).
|
||||
|
||||
func TestParseConfigDefaults(t *testing.T) {
|
||||
raw := []byte(`{
|
||||
"hosts": [{"alias":"pve","url":"https://10.0.0.5:8006/api2/json","token_id":"u@pve!mcp","token_secret":"S"}]
|
||||
}`)
|
||||
cfg, err := ParseConfig(raw)
|
||||
if err != nil {
|
||||
t.Fatalf("ParseConfig: %v", err)
|
||||
}
|
||||
if !cfg.IsReadOnly() {
|
||||
t.Error("default ReadOnly should be true")
|
||||
}
|
||||
if cfg.Default != "pve" {
|
||||
t.Errorf("default host = %q, want pve", cfg.Default)
|
||||
}
|
||||
if !cfg.HostExists("pve") {
|
||||
t.Error("pve should exist")
|
||||
}
|
||||
if cfg.MultiHost() {
|
||||
t.Error("single host should not be multi")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseConfigNoHostsFails(t *testing.T) {
|
||||
if _, err := ParseConfig([]byte(`{}`)); err == nil {
|
||||
t.Fatal("expected error for config without hosts")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseConfigMissingTokenFails(t *testing.T) {
|
||||
raw := []byte(`{
|
||||
"hosts":[{"alias":"pve","url":"https://10.0.0.5:8006/api2/json"}]
|
||||
}`)
|
||||
if _, err := ParseConfig(raw); err == nil {
|
||||
t.Fatal("expected fail-closed on missing token")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseConfigUnexpandedVarFails(t *testing.T) {
|
||||
// ${VAR} отсутствует в окружении => token_secret пустой => fail-closed.
|
||||
raw := []byte(`{
|
||||
"hosts":[{"alias":"pve","url":"https://10.0.0.5:8006/api2/json","token_id":"m","token_secret":"${DEFINITELY_MISSING_VAR}"}]
|
||||
}`)
|
||||
if _, err := ParseConfig(raw); err == nil {
|
||||
t.Fatal("expected fail-closed on unexpanded secret var")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseConfigDuplicateAliasFails(t *testing.T) {
|
||||
raw := []byte(`{
|
||||
"hosts":[
|
||||
{"alias":"pve","url":"https://10.0.0.1:8006/api2/json","token_id":"a","token_secret":"s"},
|
||||
{"alias":"pve","url":"https://10.0.0.2:8006/api2/json","token_id":"b","token_secret":"t"}
|
||||
]}`)
|
||||
if _, err := ParseConfig(raw); err == nil {
|
||||
t.Fatal("expected fail on duplicate alias")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseConfigDuplicateURLFails(t *testing.T) {
|
||||
raw := []byte(`{
|
||||
"hosts":[
|
||||
{"alias":"a","url":"https://10.0.0.1:8006/api2/json","token_id":"a","token_secret":"s"},
|
||||
{"alias":"b","url":"https://10.0.0.1:8006/api2/json","token_id":"b","token_secret":"t"}
|
||||
]}`)
|
||||
if _, err := ParseConfig(raw); err == nil {
|
||||
t.Fatal("expected fail on duplicate url")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseConfigMultiHostGlobalAllowlistFails(t *testing.T) {
|
||||
// Мульти-гипервизор + глобальный allowlist = коллизия VMID → fail-closed.
|
||||
raw := []byte(`{
|
||||
"hosts":[
|
||||
{"alias":"a","url":"https://10.0.0.1:8006/api2/json","token_id":"a","token_secret":"s"},
|
||||
{"alias":"b","url":"https://10.0.0.2:8006/api2/json","token_id":"b","token_secret":"t"}
|
||||
],
|
||||
"allowlist":{"vmids":[100]}}`)
|
||||
if _, err := ParseConfig(raw); err == nil {
|
||||
t.Fatal("expected fail on global allowlist in multi-host config")
|
||||
}
|
||||
}
|
||||
|
||||
// Ключевой тест коллизии VMID: один и тот же vmid на разных хостах
|
||||
// обязан давать РАЗНЫЙ вердикт по авторизации.
|
||||
func TestWriteAllowed_NoCrossHostLeak(t *testing.T) {
|
||||
cfg := &Config{
|
||||
Hosts: []HostConfig{
|
||||
{Alias: "a", URL: "https://10.0.0.1:8006/api2/json", TokenID: "u", TokenSecret: "s", AllowVMIDs: []int{500}},
|
||||
{Alias: "b", URL: "https://10.0.0.2:8006/api2/json", TokenID: "u", TokenSecret: "s", AllowVMIDs: []int{200}},
|
||||
},
|
||||
ReadOnly: boolPtr(false),
|
||||
}
|
||||
if !cfg.WriteAllowed("a", 500) {
|
||||
t.Error("host a vmid 500 should be allowed")
|
||||
}
|
||||
if cfg.WriteAllowed("b", 500) {
|
||||
t.Error("host b vmid 500 must be DENIED (not in its allowlist) — no cross-host leak")
|
||||
}
|
||||
if !cfg.WriteAllowed("b", 200) {
|
||||
t.Error("host b vmid 200 should be allowed")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriteAllowedFailClosedEmpty(t *testing.T) {
|
||||
cfg := &Config{Hosts: []HostConfig{{Alias: "a", URL: "https://10.0.0.1:8006/api2/json", TokenID: "u", TokenSecret: "s"}}, ReadOnly: boolPtr(false)}
|
||||
if cfg.WriteAllowed("a", 500) {
|
||||
t.Error("empty allow_vmids must fail-closed (deny writes)")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriteAllowedReadOnlyDenies(t *testing.T) {
|
||||
cfg := &Config{Hosts: []HostConfig{{Alias: "a", URL: "https://10.0.0.1:8006/api2/json", TokenID: "u", TokenSecret: "s", AllowVMIDs: []int{500}}}, ReadOnly: boolPtr(true)}
|
||||
if cfg.WriteAllowed("a", 500) {
|
||||
t.Error("read_only must deny writes even if allowlist set")
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriteAllowedGlobalFallbackSingleHost(t *testing.T) {
|
||||
// Одиночный гипервизор: глобальный allowlist — допустимый фолбэк.
|
||||
cfg := &Config{
|
||||
Hosts: []HostConfig{{Alias: "a", URL: "https://10.0.0.1:8006/api2/json", TokenID: "u", TokenSecret: "s"}},
|
||||
Allowlist: Allowlist{VMIDs: []int{500}},
|
||||
ReadOnly: boolPtr(false),
|
||||
}
|
||||
if !cfg.WriteAllowed("a", 500) {
|
||||
t.Error("global allowlist should fall back for single host")
|
||||
}
|
||||
if cfg.WriteAllowed("a", 100) {
|
||||
t.Error("vmid 100 should be denied")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNodeWriteAllowed(t *testing.T) {
|
||||
cfg := &Config{
|
||||
Hosts: []HostConfig{{Alias: "a", URL: "https://10.0.0.1:8006/api2/json", TokenID: "u", TokenSecret: "s", AllowNodes: []string{"pve"}}},
|
||||
ReadOnly: boolPtr(false),
|
||||
}
|
||||
if !cfg.NodeWriteAllowed("a", "pve") {
|
||||
t.Error("node pve should be allowed")
|
||||
}
|
||||
if cfg.NodeWriteAllowed("a", "other") {
|
||||
t.Error("node other should be denied")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDenyConfigKeys(t *testing.T) {
|
||||
cfg := &Config{denyConfigKeys: DefaultDenyConfigKeys}
|
||||
for _, k := range []string{"delete", "revert", "hotplug"} {
|
||||
if !cfg.DenyConfigKey(k) {
|
||||
t.Errorf("key %q should be denied", k)
|
||||
}
|
||||
}
|
||||
if cfg.DenyConfigKey("name") {
|
||||
t.Error("name should NOT be denied")
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsValidHTTPURL(t *testing.T) {
|
||||
valid := []string{"https://pve.local:8006/api2/json", "http://10.0.0.1:8006/api2/json"}
|
||||
for _, u := range valid {
|
||||
if !isValidHTTPURL(u) {
|
||||
t.Errorf("expected valid: %s", u)
|
||||
}
|
||||
}
|
||||
invalid := []string{"", "ftp://x", "https://", "javascript:alert(1)"}
|
||||
for _, u := range invalid {
|
||||
if isValidHTTPURL(u) {
|
||||
t.Errorf("expected invalid: %s", u)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateIdentifier(t *testing.T) {
|
||||
ok := []string{"pve", "pve1", "snap-name", "104", "local-lvm", "scsi0"}
|
||||
for _, s := range ok {
|
||||
if err := ValidateIdentifier(s); err != nil {
|
||||
t.Errorf("expected valid %q: %v", s, err)
|
||||
}
|
||||
}
|
||||
bad := []string{"", "a/b", "..", "a b", "a;rm", "a\\b", "a\x00b", "a:b", "a@b"}
|
||||
for _, s := range bad {
|
||||
if err := ValidateIdentifier(s); err == nil {
|
||||
t.Errorf("expected invalid %q", s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateUPID(t *testing.T) {
|
||||
ok := "UPID:pve:00000000:root@pam!mcp:1:2:3:qemu:100:abc"
|
||||
if err := ValidateUPID(ok); err != nil {
|
||||
t.Errorf("expected valid upid %q: %v", ok, err)
|
||||
}
|
||||
bad := []string{"", "a/b", "..", "a\\b", "a\x00b"}
|
||||
for _, s := range bad {
|
||||
if err := ValidateUPID(s); err == nil {
|
||||
t.Errorf("expected invalid upid %q", s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func boolPtr(b bool) *bool { return &b }
|
||||
@@ -0,0 +1,154 @@
|
||||
package pve
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"sync"
|
||||
|
||||
"git.totmin.ru/en2zmax/forge-toolkit/configreload"
|
||||
)
|
||||
|
||||
// Manager — одиночный, ПОТОКО-БЕЗОПАСНЫЙ диспетчер процесса (столп 2 §2):
|
||||
// агенты запускаются в отдельных goroutine и конкурентно дёргают один
|
||||
// Manager. В pooled-режиме (isolation=pooled) ядро на каждый вызов
|
||||
// инжектит серверный аргумент _tenant_config = <agentDir>/forge-tools/
|
||||
// proxmox.json — поэтому менеджер кэширует per-путь загрузчики конфига
|
||||
// (configreload, live-reload по контент-хэшу) и per-путь Tenant'ы (свои
|
||||
// клиенты/секреты). Тенанты НЕ делят клиентов между агентами.
|
||||
type Manager struct {
|
||||
mu sync.Mutex
|
||||
// configPath — статический -config (легаси-одиночный режим). Пусто =
|
||||
// pooled: конфиг приходит на каждый вызов через _tenant_config.
|
||||
configPath string
|
||||
// loaders кэширует configreload.Loader[*Config] по пути конфига.
|
||||
loaders map[string]*configreload.Loader[*Config]
|
||||
// tenants кэширует Tenant (свои клиенты) по ключу-пути конфига.
|
||||
tenants map[string]*Tenant
|
||||
}
|
||||
|
||||
// NewManager создаёт менеджер. configPath — опциональный постоянный конфиг
|
||||
// (-config); пусто = pooled (тенант из _tenant_config на каждый вызов).
|
||||
func NewManager(configPath string) *Manager {
|
||||
return &Manager{
|
||||
configPath: configPath,
|
||||
loaders: make(map[string]*configreload.Loader[*Config]),
|
||||
tenants: make(map[string]*Tenant),
|
||||
}
|
||||
}
|
||||
|
||||
// Tenant возвращает per-agent тенант. path — _tenant_config (пусто при
|
||||
// -config-режиме). Fail-closed: нет конфига — ошибка (модуль отказывает,
|
||||
// а не работает «с общими» кредами).
|
||||
func (m *Manager) Tenant(ctx context.Context, path string) (*Tenant, error) {
|
||||
cfgPath := m.resolvePath(path)
|
||||
if cfgPath == "" {
|
||||
return nil, errors.New("proxmox: no config (set -config or provide _tenant_config)")
|
||||
}
|
||||
|
||||
loader := m.loaderFor(cfgPath)
|
||||
cfg, err := loader.Get()
|
||||
if err != nil {
|
||||
// ErrNotFound / parse-error без last-good — fail-closed, а не фолбэк.
|
||||
if cfg == nil || !cfg.HasHosts() {
|
||||
return nil, fmt.Errorf("proxmox: tenant config %s: %w", cfgPath, err)
|
||||
}
|
||||
// Есть last-good — работаем со старым (правка была битой), но
|
||||
// сигналим, чтобы не молчать.
|
||||
// (здесь последний рабочий конфиг уже возвращён в cfg)
|
||||
}
|
||||
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if t, ok := m.tenants[cfgPath]; ok {
|
||||
return t, nil
|
||||
}
|
||||
t := &Tenant{cfg: cfg, clients: make(map[string]*Client)}
|
||||
m.tenants[cfgPath] = t
|
||||
return t, nil
|
||||
}
|
||||
|
||||
// resolvePath выбирает путь конфига: _tenant_config в приоритете, иначе
|
||||
// статический -config.
|
||||
func (m *Manager) resolvePath(tenantPath string) string {
|
||||
if tenantPath != "" {
|
||||
return tenantPath
|
||||
}
|
||||
return m.configPath
|
||||
}
|
||||
|
||||
// loaderFor возвращает (и кэширует) загрузчик конфига по пути.
|
||||
func (m *Manager) loaderFor(path string) *configreload.Loader[*Config] {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if l, ok := m.loaders[path]; ok {
|
||||
return l
|
||||
}
|
||||
l := configreload.New(path, ParseConfig)
|
||||
m.loaders[path] = l
|
||||
return l
|
||||
}
|
||||
|
||||
// Close закрывает все тенанты (и их клиенты).
|
||||
func (m *Manager) Close() {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
for _, t := range m.tenants {
|
||||
t.closeLocked()
|
||||
}
|
||||
m.tenants = nil
|
||||
}
|
||||
|
||||
// Tenant — per-agent конфигурация + свои HTTP-клиенты к хостам. Секреты и
|
||||
// разрешения (allowlist) — строго в рамках одного тенанта.
|
||||
type Tenant struct {
|
||||
cfg *Config
|
||||
mu sync.Mutex
|
||||
clients map[string]*Client // alias -> client
|
||||
}
|
||||
|
||||
// Config возвращает политику тенанта.
|
||||
func (t *Tenant) Config() *Config { return t.cfg }
|
||||
|
||||
// Client возвращает HTTP-клиент для хоста по алиасу (или default).
|
||||
// Строится лениво и кэшируется; потокобезопасно.
|
||||
func (t *Tenant) Client(ctx context.Context, alias string) (*Client, error) {
|
||||
host, ok := t.cfg.Host(alias)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("proxmox: host %q not configured", alias)
|
||||
}
|
||||
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
if c, ok := t.clients[host.Alias]; ok {
|
||||
return c, nil
|
||||
}
|
||||
c, err := NewClient(host)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
t.clients[host.Alias] = c
|
||||
return c, nil
|
||||
}
|
||||
|
||||
// Close закрывает клиенты тенанта. Не используй вне RWMutex Manager.
|
||||
func (t *Tenant) Close() {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
t.closeLocked()
|
||||
}
|
||||
|
||||
func (t *Tenant) closeLocked() {
|
||||
// net/http.Client не имеет Close; здесь точка для будущего пула/окружения.
|
||||
t.clients = nil
|
||||
}
|
||||
|
||||
// joinURL собирает корректный путь API (PathEscape против path-traversal).
|
||||
func joinURL(path string, ids ...string) string {
|
||||
p := path
|
||||
for _, id := range ids {
|
||||
p += "/" + url.PathEscape(id)
|
||||
}
|
||||
return p
|
||||
}
|
||||
Reference in New Issue
Block a user