// 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) } // 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"` 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 }