Files
felhom.eu/hub/internal/tenantsync/client.go
T
admin 4009401f46 hub v0.61.0 + felhom-tenantsync v1.1.0: Customer RESET (middle lifecycle tier)
One operator action returns a customer to pre-first-install: all operational
state dies (offsite repo, PBS namespace+backups, DR recipe, one-time secret,
claim state, retained escrow custody); identity + basic config + provenance +
events survive. Sits between host delete and customer Delete.

- store/customer_reset.go: customer_resets journal, live inventory, ack-gated
  purge (never touches identity/provenance/events), DeleteClaim.
- claim.ResetToUnclaimed: delete claim row -> fresh code next onboarding.
- offsite.Deprovision (idempotent) + OffsiteIdentifier + ClearProvisionedDescriptor.
- tenantsync.Deprovision + felhom-tenantsync.sh deprovision op (destroys ns +
  backup groups + token; shared user untouched; idempotent).
- web/customer_reset.go: GET reset -> inventory JSON; POST -> orchestration
  (external teardown FIRST, DB purge LAST; refuse-while-hosts; typed-id +
  separate escrow ack). Amber RESET card distinct from red Danger-zone Delete.
- Red-proofs: ack-gate + partial-failure resumability (both proven red);
  store ack-gating + journal round-trip; offsite idempotency + descriptor clear;
  RESET-card render. Green: build + vet + test.
2026-07-17 13:09:04 +02:00

246 lines
9.7 KiB
Go

// Package tenantsync drives the offsite endpoint's per-customer PBS tenancy surface (PBS DR tier
// SLICE 1) — the structural twin of internal/wgsync: SSH with a PINNED host key (exact-match or
// refuse, no fallback) to a forced-command script (`felhom-tenantsync` — its OWN key + sudoers
// line; the peersync surface is untouched). JSON on stdin, JSON on stdout, one op per session.
//
// SECRET HYGIENE (load-bearing divergence from wgsync): the script's stdout carries the one-time
// PBS token secret. It is parsed into Result and handed to the caller for the consume-once store
// write — it is NEVER logged, and error messages NEVER embed stdout bytes (stderr only). Do not
// "improve" the diagnostics by quoting the response.
package tenantsync
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"log"
"net"
"regexp"
"strings"
"time"
"golang.org/x/crypto/ssh"
)
// ErrTokenExists is the typed "provision refused: the token already exists" outcome — the hub
// treats it as a state mismatch (a descriptor should exist; re-issue is the explicit recovery).
var ErrTokenExists = errors.New("tenantsync: token already exists on the endpoint (re-issue is the explicit path)")
// customerIDRe mirrors the script's validation — refuse client-side before a wasted SSH round-trip.
var customerIDRe = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9_.-]{0,30}$`)
// Config configures the SSH client. Addr/User/HostKeyLine are typically the SAME values as the
// peersync client (same box, same low-priv user, same pinned host key); PrivateKey is tenantsync's
// OWN key (the authorized_keys line selects the forced command).
type Config struct {
Addr string // "host:22"
User string // "felhom-peersync"
PrivateKey []byte // PEM private key (from the mounted Secret file)
HostKeyLine string // single authorized_keys-format line of the endpoint's host pubkey
Timeout time.Duration // default 60s (tenancy ops run several PBS commands)
}
// Result is the script's ok-response: the descriptor fields + the ONE-TIME token secret.
// TokenSecret is transient custody — store it consume-once immediately, never log the struct.
type Result struct {
TokenID string `json:"token_id"`
TokenSecret string `json:"token_secret"`
Fingerprint string `json:"fingerprint"`
Datastore string `json:"datastore"`
Namespace string `json:"namespace"`
}
// Client is a pinned-host-key SSH per-op executor. Construct with New (parses keys up front).
type Client struct {
addr string
user string
signer ssh.Signer
hostKey ssh.PublicKey
timeout time.Duration
logger *log.Logger
}
// New builds a Client, failing early on an unparsable private key or host-key line.
func New(cfg Config, logger *log.Logger) (*Client, error) {
if cfg.Addr == "" || cfg.User == "" {
return nil, fmt.Errorf("tenantsync: Addr and User are required")
}
signer, err := ssh.ParsePrivateKey(cfg.PrivateKey)
if err != nil {
return nil, fmt.Errorf("tenantsync: parse private key: %w", err)
}
hostKey, _, _, _, err := ssh.ParseAuthorizedKey([]byte(cfg.HostKeyLine))
if err != nil {
return nil, fmt.Errorf("tenantsync: parse host key line: %w", err)
}
timeout := cfg.Timeout
if timeout == 0 {
timeout = 60 * time.Second
}
if logger == nil {
logger = log.Default()
}
return &Client{
addr: cfg.Addr, user: cfg.User, signer: signer,
hostKey: hostKey, timeout: timeout, logger: logger,
}, nil
}
// Provision creates the customer's namespace + privilege-separated token on the endpoint.
// An already-existing token is the typed ErrTokenExists (never silently re-keyed).
func (c *Client) Provision(ctx context.Context, customerID string) (*Result, error) {
return c.tenancyOp(ctx, "provision", customerID)
}
// Reissue explicitly re-keys the customer's token (delete + recreate + re-grant, script-side).
func (c *Client) Reissue(ctx context.Context, customerID string) (*Result, error) {
return c.tenancyOp(ctx, "reissue", customerID)
}
// Deprovision DESTROYS the customer's PBS namespace, all its backup groups, and its token — the
// customer-RESET teardown (v0.61.0, operator ack-gated). The shared felhom@pbs user is never touched
// (co-tenants ride it). Idempotent: a missing tenant is a clean success (existed=false). This carries
// NO secret, so it does not route through tenancyOp's token-field validation.
func (c *Client) Deprovision(ctx context.Context, customerID string) (existed bool, err error) {
if !customerIDRe.MatchString(customerID) {
return false, fmt.Errorf("tenantsync: invalid customer_id %q", customerID)
}
payload, err := json.Marshal(map[string]string{"op": "deprovision", "customer_id": customerID})
if err != nil {
return false, err
}
stdout, stderr, runErr := c.exec(ctx, payload)
resp, err := parseResponse(stdout, stderr, runErr)
if err != nil {
return false, err
}
if resp.Namespace == "" {
return false, fmt.Errorf("tenantsync: deprovision response missing namespace")
}
c.logger.Printf("[INFO] tenantsync: deprovision ok for %s (ns=%s, existed=%t)",
customerID, resp.Namespace, resp.Deleted)
return resp.Deleted, nil
}
// Fingerprint returns the endpoint PBS's API cert fingerprint (the descriptor field).
func (c *Client) Fingerprint(ctx context.Context) (string, error) {
stdout, stderr, runErr := c.exec(ctx, []byte(`{"op":"fingerprint"}`))
resp, err := parseResponse(stdout, stderr, runErr)
if err != nil {
return "", err
}
if resp.Fingerprint == "" {
return "", fmt.Errorf("tenantsync: fingerprint op returned an empty fingerprint")
}
return resp.Fingerprint, nil
}
func (c *Client) tenancyOp(ctx context.Context, op, customerID string) (*Result, error) {
if !customerIDRe.MatchString(customerID) {
return nil, fmt.Errorf("tenantsync: invalid customer_id %q", customerID)
}
payload, err := json.Marshal(map[string]string{"op": op, "customer_id": customerID})
if err != nil {
return nil, err
}
stdout, stderr, runErr := c.exec(ctx, payload)
resp, err := parseResponse(stdout, stderr, runErr)
if err != nil {
return nil, err
}
if resp.TokenSecret == "" || resp.TokenID == "" || resp.Namespace == "" || resp.Fingerprint == "" || resp.Datastore == "" {
// Field NAMES only — never values (TokenSecret).
return nil, fmt.Errorf("tenantsync: %s response is missing required fields", op)
}
c.logger.Printf("[INFO] tenantsync: %s ok for %s (ns=%s, token_id=%s; secret withheld from logs)",
op, customerID, resp.Namespace, resp.TokenID)
return &resp.Result, nil
}
// response is the script's stdout contract — ok carries the Result fields, error carries code+error.
type response struct {
Status string `json:"status"`
Code string `json:"code"`
Error string `json:"error"`
Deleted bool `json:"deleted"` // deprovision op: whether the namespace existed (was destroyed)
Result
}
// parseResponse turns (stdout, stderr, runErr) into a typed outcome. The script emits its error
// JSON on stdout and exits 1, so a run error is parsed for the typed code FIRST; only when stdout
// carries no usable JSON does the raw failure (with stderr, NEVER stdout) surface.
func parseResponse(stdout, stderr []byte, runErr error) (*response, error) {
var resp response
parseOK := json.Unmarshal(bytes.TrimSpace(stdout), &resp) == nil
if parseOK && resp.Status == "error" {
if resp.Code == "token_exists" {
return nil, ErrTokenExists
}
return nil, fmt.Errorf("tenantsync: endpoint refused: %s (code %s)", resp.Error, resp.Code)
}
if runErr != nil {
return nil, fmt.Errorf("tenantsync: remote op failed: %w (stderr: %s)",
runErr, strings.TrimSpace(string(stderr)))
}
if !parseOK || resp.Status != "ok" {
// stdout may carry the secret — report shape only, never bytes.
return nil, fmt.Errorf("tenantsync: malformed endpoint response (%d stdout bytes; stderr: %s)",
len(bytes.TrimSpace(stdout)), strings.TrimSpace(string(stderr)))
}
return &resp, nil
}
// exec runs one forced-command session: payload on stdin, returns stdout/stderr. The connection
// discipline (pinned key, constrained HostKeyAlgorithms, deadline both sides of the handshake)
// is wgsync.Push's, verbatim — that shape is live-proven against the real sshd.
func (c *Client) exec(ctx context.Context, payload []byte) (stdout, stderr []byte, err error) {
sshCfg := &ssh.ClientConfig{
User: c.user,
Auth: []ssh.AuthMethod{ssh.PublicKeys(c.signer)},
HostKeyCallback: ssh.FixedHostKey(c.hostKey),
// Constrain negotiation to the PINNED key's algorithm — a multi-hostkey sshd (stock:
// ECDSA + ed25519) otherwise presents a different type and FixedHostKey refuses a
// legitimate server (wgsync S1 live finding).
HostKeyAlgorithms: []string{c.hostKey.Type()},
Timeout: c.timeout,
}
dialer := net.Dialer{Timeout: c.timeout}
conn, err := dialer.DialContext(ctx, "tcp", c.addr)
if err != nil {
return nil, nil, fmt.Errorf("tenantsync: dial %s: %w", c.addr, err)
}
if dl, ok := ctx.Deadline(); ok {
conn.SetDeadline(dl)
} else {
conn.SetDeadline(time.Now().Add(c.timeout))
}
sconn, chans, reqs, err := ssh.NewClientConn(conn, c.addr, sshCfg)
if err != nil {
conn.Close()
return nil, nil, fmt.Errorf("tenantsync: ssh handshake %s: %w", c.addr, err)
}
client := ssh.NewClient(sconn, chans, reqs)
defer client.Close()
conn.SetDeadline(time.Time{})
if dl, ok := ctx.Deadline(); ok {
conn.SetDeadline(dl)
}
session, err := client.NewSession()
if err != nil {
return nil, nil, fmt.Errorf("tenantsync: session: %w", err)
}
defer session.Close()
var outBuf, errBuf bytes.Buffer
session.Stdin = bytes.NewReader(payload)
session.Stdout = &outBuf
session.Stderr = &errBuf
// The forced command overrides this string; it documents intent on the wire.
runErr := session.Run("felhom-tenantsync")
return outBuf.Bytes(), errBuf.Bytes(), runErr
}