feat(remote): route LXC/VM checks through the Proxmox host, not direct SSH
The scheduler SSHed each guest directly and assumed a deployed probe script plus working root SSH at the guest's address — false for headless (nfs-export), keyless (teddycloud), mesh-only (rclone), and macOS (mac-mini) targets, which left 49 enabled checks stuck "down" on a healthy fleet. Extract the MCP run tool's resolveExecTarget into a shared internal/remote package and make it the single execution path for both the scheduler and MCP. LXC/VM checks now host-hop via pct exec / qm guest exec through the owning Proxmox host (no per-guest lan_ip, sshd, or authorized key needed); hosts and workstations resolve their address and user live, so mac-mini's `user: dtoro` is honored without a re-seed. Address preference now prefers public_ipv4 over mesh, so netbird-vps is probeable from the scheduler container. cpu_check.sh gains a real Darwin branch (it reported cpu_pct 0 before). checkdefaults.resolveSSHUser reads the top-level `user` attribute too. A machine-target resolution failure is now logged before falling back to baked config, so a broken probe-config is distinguishable from a real outage.
This commit is contained in:
257
internal/remote/remote.go
Normal file
257
internal/remote/remote.go
Normal file
@@ -0,0 +1,257 @@
|
||||
// Package remote resolves how to execute a command on a target entity and
|
||||
// turns a plain shell command into whatever must be sent over the SSH
|
||||
// connection that reaches it.
|
||||
//
|
||||
// The canonical access model: a host or workstation is reached by direct SSH
|
||||
// to its address; an LXC or VM is NEVER SSH'd into directly — it is reached
|
||||
// through its owning Proxmox host via `pct exec` / `qm guest exec`. One SSH
|
||||
// credential per host (the host's root key), no per-guest keys, sshd, or
|
||||
// lan_ip required for execution. Network probes (http/ping) still hit a
|
||||
// guest's lan_ip directly; only command execution host-hops.
|
||||
//
|
||||
// This is the single resolver shared by the scheduler's check execution and
|
||||
// the MCP `run` tool. Previously they diverged — the scheduler SSHed guests
|
||||
// directly (broken for headless/keyless/mesh-only guests), while MCP
|
||||
// host-hopped (working). Keeping one path keeps them in lockstep.
|
||||
package remote
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/dtoro/oikos/internal/db"
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
// DefaultUser is the SSH user when an entity declares no ssh.user. The
|
||||
// Proxmox hosts and their guests are all administered as root.
|
||||
const DefaultUser = "root"
|
||||
|
||||
// ExecTarget is a resolved execution endpoint: the SSH address and user to
|
||||
// connect to, plus Wrap, which rewrites a plain command for transport.
|
||||
type ExecTarget struct {
|
||||
Host string
|
||||
User string
|
||||
// Wrap turns a plain shell command into the form that must be sent over
|
||||
// the SSH connection to this target: the identity function for a host,
|
||||
// `pct exec <id> -- bash -c 'echo <b64> | base64 -d | bash'` for an LXC,
|
||||
// the `qm guest exec` equivalent for a VM. The base64 round-trip keeps
|
||||
// nested quoting identical across both guest kinds.
|
||||
Wrap func(cmd string) string
|
||||
}
|
||||
|
||||
// IsGuest reports whether an entity type is reached via pct/qm exec through a
|
||||
// Proxmox host rather than by direct SSH. docker-container is reached via its
|
||||
// host's docker socket, not pct, so it is not a guest here.
|
||||
func IsGuest(entityType string) bool {
|
||||
return entityType == "lxc" || entityType == "vm"
|
||||
}
|
||||
|
||||
// ResolveHost resolves a `host:<slug>` to its reachable network address and
|
||||
// SSH user. Address preference: lan_ip, then public_ipv4, then mesh IP, then
|
||||
// mesh fqdn. Preferring public_ipv4 over mesh matters because the scheduler
|
||||
// container has no mesh interface — a standalone-server with only a mesh IP
|
||||
// (netbird-vps) was unreachable, and a public_ipv4 was sitting unused.
|
||||
//
|
||||
// fallbackUser is used when the entity declares no ssh.user; callers pass
|
||||
// their configured default (the scheduler uses "root", the MCP run tool uses
|
||||
// its configured OIKOS_SSH_USER).
|
||||
func ResolveHost(ctx context.Context, pool *db.Pool, hostSlug, fallbackUser string) (addr, user string, err error) {
|
||||
var raw string
|
||||
if err = pool.QueryRow(ctx, "SELECT attributes::text FROM entities WHERE slug = $1", hostSlug).Scan(&raw); err != nil {
|
||||
return "", "", fmt.Errorf("entity not found: %s", hostSlug)
|
||||
}
|
||||
var m map[string]any
|
||||
if err = json.Unmarshal([]byte(raw), &m); err != nil {
|
||||
return "", "", fmt.Errorf("parse attributes for %s: %w", hostSlug, err)
|
||||
}
|
||||
|
||||
if v, ok := m["lan_ip"].(string); ok && v != "" {
|
||||
addr = v
|
||||
} else if v, ok := m["public_ipv4"].(string); ok && v != "" {
|
||||
addr = v
|
||||
} else if mesh, ok := m["mesh"].(map[string]any); ok {
|
||||
if nb, ok := mesh["netbird"].(map[string]any); ok {
|
||||
if v, ok := nb["ip"].(string); ok && v != "" {
|
||||
addr = v
|
||||
} else if v, ok := nb["fqdn"].(string); ok && v != "" {
|
||||
addr = v
|
||||
}
|
||||
}
|
||||
}
|
||||
if addr == "" {
|
||||
return "", "", fmt.Errorf("no IP found for %s", hostSlug)
|
||||
}
|
||||
|
||||
user = fallbackUser
|
||||
if ssh, ok := m["ssh"].(map[string]any); ok {
|
||||
if u, ok := ssh["user"].(string); ok && u != "" {
|
||||
user = u
|
||||
}
|
||||
}
|
||||
return addr, user, nil
|
||||
}
|
||||
|
||||
// ResolveProxmoxHostSlug resolves the Proxmox host slug that owns a guest.
|
||||
// Resolution order: the hostAttr if non-empty (the entity's attributes.host,
|
||||
// stored without the "host:" prefix), the `hosts` relationship on the guest
|
||||
// (the canonical graph edge), then "hubris" as the documented default.
|
||||
//
|
||||
// entityID is the guest's entity id; the relationship lookup uses it
|
||||
// directly rather than a slug subquery.
|
||||
func ResolveProxmoxHostSlug(ctx context.Context, pool *db.Pool, entityID uuid.UUID, hostAttr string) string {
|
||||
hostSlug := strings.TrimSpace(hostAttr)
|
||||
if hostSlug == "" {
|
||||
var relHostSlug string
|
||||
if err := pool.QueryRow(ctx, `
|
||||
SELECT e.slug FROM relationships r
|
||||
JOIN entities e ON e.id = r.source_id
|
||||
WHERE r.target_id = $1
|
||||
AND r.type = 'hosts' AND r.valid_to IS NULL
|
||||
LIMIT 1`, entityID).Scan(&relHostSlug); err == nil && relHostSlug != "" {
|
||||
hostSlug = relHostSlug
|
||||
}
|
||||
}
|
||||
if hostSlug == "" {
|
||||
hostSlug = "hubris" // documented default Proxmox host when unset
|
||||
}
|
||||
if !strings.HasPrefix(hostSlug, "host:") {
|
||||
hostSlug = "host:" + hostSlug
|
||||
}
|
||||
return hostSlug
|
||||
}
|
||||
|
||||
// ResolveExecTarget resolves a target slug (host:, lxc:, or vm:) to its
|
||||
// execution endpoint. This is the slug-based entry used by the MCP `run` tool.
|
||||
func ResolveExecTarget(ctx context.Context, pool *db.Pool, targetSlug, fallbackUser string) (ExecTarget, error) {
|
||||
switch {
|
||||
case strings.HasPrefix(targetSlug, "host:"):
|
||||
addr, user, err := ResolveHost(ctx, pool, targetSlug, fallbackUser)
|
||||
if err != nil {
|
||||
return ExecTarget{}, err
|
||||
}
|
||||
return ExecTarget{Host: addr, User: user, Wrap: func(cmd string) string { return cmd }}, nil
|
||||
|
||||
case strings.HasPrefix(targetSlug, "lxc:"), strings.HasPrefix(targetSlug, "vm:"):
|
||||
var (
|
||||
id uuid.UUID
|
||||
pveID string
|
||||
typ string
|
||||
hostAttr string
|
||||
)
|
||||
if err := pool.QueryRow(ctx,
|
||||
"SELECT id, type, attributes->>'pve_id', COALESCE(attributes->>'host','') FROM entities WHERE slug = $1",
|
||||
targetSlug).Scan(&id, &typ, &pveID, &hostAttr); err != nil || pveID == "" {
|
||||
return ExecTarget{}, fmt.Errorf("guest not found or missing pve_id: %s", targetSlug)
|
||||
}
|
||||
hostSlug := ResolveProxmoxHostSlug(ctx, pool, id, hostAttr)
|
||||
addr, user, err := ResolveHost(ctx, pool, hostSlug, fallbackUser)
|
||||
if err != nil {
|
||||
return ExecTarget{}, err
|
||||
}
|
||||
return ExecTarget{Host: addr, User: user, Wrap: guestWrap(typ, pveID)}, nil
|
||||
}
|
||||
return ExecTarget{}, fmt.Errorf("unsupported target %q: must be host:<slug>, lxc:<slug>, or vm:<slug>", targetSlug)
|
||||
}
|
||||
|
||||
// ResolveExecTargetForCheck resolves an execution endpoint keyed by the
|
||||
// target's id and type — the data the scheduler has at check-execution time
|
||||
// (check_defs carry target_id + target_type, not a slug). Guests route via
|
||||
// pct/qm exec; everything else (hosts, workstations, services resolved to
|
||||
// their hosting machine) is reached by direct SSH to the entity's own address.
|
||||
func ResolveExecTargetForCheck(ctx context.Context, pool *db.Pool, targetID uuid.UUID, targetType, fallbackUser string) (ExecTarget, error) {
|
||||
if IsGuest(targetType) {
|
||||
var (
|
||||
pveID string
|
||||
hostAttr string
|
||||
)
|
||||
if err := pool.QueryRow(ctx,
|
||||
"SELECT attributes->>'pve_id', COALESCE(attributes->>'host','') FROM entities WHERE id = $1",
|
||||
targetID).Scan(&pveID, &hostAttr); err != nil || pveID == "" {
|
||||
return ExecTarget{}, fmt.Errorf("guest %s missing pve_id", targetID)
|
||||
}
|
||||
hostSlug := ResolveProxmoxHostSlug(ctx, pool, targetID, hostAttr)
|
||||
addr, user, err := ResolveHost(ctx, pool, hostSlug, fallbackUser)
|
||||
if err != nil {
|
||||
return ExecTarget{}, err
|
||||
}
|
||||
return ExecTarget{Host: addr, User: user, Wrap: guestWrap(targetType, pveID)}, nil
|
||||
}
|
||||
|
||||
// Host-like target: reach it directly at its own address. Services and
|
||||
// other non-host entities that reach here should already have had their
|
||||
// host address baked into check config at seed time; this path covers
|
||||
// host/workstation targets whose address is resolved live.
|
||||
addr, user, err := resolveHostByID(ctx, pool, targetID, fallbackUser)
|
||||
if err != nil {
|
||||
return ExecTarget{}, err
|
||||
}
|
||||
return ExecTarget{Host: addr, User: user, Wrap: func(cmd string) string { return cmd }}, nil
|
||||
}
|
||||
|
||||
// resolveHostByID is ResolveHost keyed by entity id.
|
||||
func resolveHostByID(ctx context.Context, pool *db.Pool, id uuid.UUID, fallbackUser string) (addr, user string, err error) {
|
||||
var raw string
|
||||
if err = pool.QueryRow(ctx, "SELECT attributes::text FROM entities WHERE id = $1", id).Scan(&raw); err != nil {
|
||||
return "", "", fmt.Errorf("entity %s not found", id)
|
||||
}
|
||||
var m map[string]any
|
||||
if err = json.Unmarshal([]byte(raw), &m); err != nil {
|
||||
return "", "", fmt.Errorf("parse attributes: %w", err)
|
||||
}
|
||||
for _, key := range []string{"lan_ip", "public_ipv4", "mesh_ip"} {
|
||||
if v, ok := m[key].(string); ok && v != "" {
|
||||
addr = v
|
||||
break
|
||||
}
|
||||
}
|
||||
if addr == "" {
|
||||
if mesh, ok := m["mesh"].(map[string]any); ok {
|
||||
if nb, ok := mesh["netbird"].(map[string]any); ok {
|
||||
if v, ok := nb["ip"].(string); ok && v != "" {
|
||||
addr = v
|
||||
} else if v, ok := nb["fqdn"].(string); ok && v != "" {
|
||||
addr = v
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if addr == "" {
|
||||
return "", "", fmt.Errorf("no IP found for entity %s", id)
|
||||
}
|
||||
user = fallbackUser
|
||||
if ssh, ok := m["ssh"].(map[string]any); ok {
|
||||
if u, ok := ssh["user"].(string); ok && u != "" {
|
||||
user = u
|
||||
}
|
||||
}
|
||||
if u, ok := m["user"].(string); ok && u != "" && user == fallbackUser {
|
||||
// Workstations carry their login as a top-level `user` attribute
|
||||
// (mac-mini: user: dtoro), not under ssh.user. Take it only when no
|
||||
// explicit ssh.user was set, so a host that genuinely wants root still
|
||||
// gets root.
|
||||
user = u
|
||||
}
|
||||
return addr, user, nil
|
||||
}
|
||||
|
||||
// guestWrap builds the pct/qm exec wrapper for a guest of the given type.
|
||||
func guestWrap(entityType, pveID string) func(cmd string) string {
|
||||
if entityType == "vm" {
|
||||
return func(cmd string) string {
|
||||
b64 := base64.StdEncoding.EncodeToString([]byte(cmd))
|
||||
// `qm guest exec` returns JSON; pipe through jq for a clean stdout,
|
||||
// falling back to the raw form. Mirrors the LXC base64 round-trip.
|
||||
return fmt.Sprintf(
|
||||
"qm guest exec %s -- /bin/bash -c 'echo %s | base64 -d | bash' | jq -r '.out // .err // empty' 2>/dev/null || qm guest exec %s -- /bin/bash -c 'echo %s | base64 -d | bash'",
|
||||
pveID, b64, pveID, b64)
|
||||
}
|
||||
}
|
||||
return func(cmd string) string {
|
||||
b64 := base64.StdEncoding.EncodeToString([]byte(cmd))
|
||||
return fmt.Sprintf("pct exec %s -- bash -c 'echo %s | base64 -d | bash'", pveID, b64)
|
||||
}
|
||||
}
|
||||
172
internal/remote/remote_test.go
Normal file
172
internal/remote/remote_test.go
Normal file
@@ -0,0 +1,172 @@
|
||||
package remote
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/dtoro/oikos/internal/db"
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
// guestWrap and IsGuest are pure logic — always tested. The DB-backed
|
||||
// resolvers are integration tests guarded by OIKOS_TEST_DATABASE_URL, the
|
||||
// same convention as internal/scheduler/coverage_test.go.
|
||||
|
||||
func TestIsGuest(t *testing.T) {
|
||||
cases := map[string]bool{
|
||||
"lxc": true, "vm": true,
|
||||
"proxmox-host": false, "workstation": false,
|
||||
"service": false, "docker-container": false,
|
||||
}
|
||||
for typ, want := range cases {
|
||||
if got := IsGuest(typ); got != want {
|
||||
t.Errorf("IsGuest(%q) = %v, want %v", typ, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestGuestWrapLXC(t *testing.T) {
|
||||
w := guestWrap("lxc", "132")
|
||||
out := w("/opt/oikos/checks/cpu_check.sh 'svc'")
|
||||
if !strings.Contains(out, "pct exec 132 -- bash -c ") {
|
||||
t.Fatalf("lxc wrap must use pct exec: %q", out)
|
||||
}
|
||||
if strings.Contains(out, "qm guest exec") {
|
||||
t.Fatalf("lxc wrap must not use qm: %q", out)
|
||||
}
|
||||
// The base64 payload must round-trip to the original command.
|
||||
i := strings.Index(out, "echo ")
|
||||
j := strings.LastIndex(out, " | base64 -d | bash")
|
||||
if i < 0 || j < 0 || j <= i {
|
||||
t.Fatalf("cannot locate base64 payload in %q", out)
|
||||
}
|
||||
dec, err := base64.StdEncoding.DecodeString(out[i+len("echo ") : j])
|
||||
if err != nil {
|
||||
t.Fatalf("decode payload: %v", err)
|
||||
}
|
||||
if string(dec) != "/opt/oikos/checks/cpu_check.sh 'svc'" {
|
||||
t.Fatalf("round-trip mismatch: %q", string(dec))
|
||||
}
|
||||
}
|
||||
|
||||
func TestGuestWrapVM(t *testing.T) {
|
||||
w := guestWrap("vm", "108")
|
||||
out := w("uname -a")
|
||||
if !strings.Contains(out, "qm guest exec 108") {
|
||||
t.Fatalf("vm wrap must use qm guest exec: %q", out)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResolveExecTargetUnsupported(t *testing.T) {
|
||||
// No DB needed: an unsupported slug prefix errors before any query.
|
||||
if _, err := ResolveExecTarget(context.Background(), nil, "service:gitea", DefaultUser); err == nil {
|
||||
t.Fatal("expected error for unsupported target prefix")
|
||||
}
|
||||
}
|
||||
|
||||
// --- integration tests (require a real Postgres) ---
|
||||
|
||||
func newRemotePool(t *testing.T) *db.Pool {
|
||||
t.Helper()
|
||||
base := testDatabaseURL(t)
|
||||
return createTestDB(t, base)
|
||||
}
|
||||
|
||||
func testDatabaseURL(t *testing.T) string {
|
||||
t.Helper()
|
||||
u := getenvOrDefault("OIKOS_TEST_DATABASE_URL", "")
|
||||
if u == "" {
|
||||
t.Skip("OIKOS_TEST_DATABASE_URL not set — skipping integration test")
|
||||
}
|
||||
return u
|
||||
}
|
||||
|
||||
func TestResolveHostPrefersLAN(t *testing.T) {
|
||||
pool := newRemotePool(t)
|
||||
ctx := context.Background()
|
||||
mustExec(t, pool, ctx, `INSERT INTO entities (id, slug, type, name, state, attributes, version, created_at, updated_at)
|
||||
VALUES ($1,'host:x','proxmox-host','x','active','{"lan_ip":"10.0.0.1","public_ipv4":"1.2.3.4","mesh":{"netbird":{"ip":"100.64.0.1"}}}'::jsonb,1,now(),now())`,
|
||||
uuid.New())
|
||||
|
||||
addr, user, err := ResolveHost(ctx, pool, "host:x", DefaultUser)
|
||||
if err != nil {
|
||||
t.Fatalf("ResolveHost: %v", err)
|
||||
}
|
||||
if addr != "10.0.0.1" {
|
||||
t.Errorf("addr = %q, want lan_ip 10.0.0.1", addr)
|
||||
}
|
||||
if user != "root" {
|
||||
t.Errorf("user = %q, want root", user)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResolveHostFallsBackToPublicIPv4(t *testing.T) {
|
||||
// netbird-vps: no lan_ip, has public_ipv4 + mesh ip. Must prefer
|
||||
// public_ipv4 — the scheduler container has no mesh interface.
|
||||
pool := newRemotePool(t)
|
||||
ctx := context.Background()
|
||||
mustExec(t, pool, ctx, `INSERT INTO entities (id, slug, type, name, state, attributes, version, created_at, updated_at)
|
||||
VALUES ($1,'host:vps','standalone-server','vps','active','{"public_ipv4":"82.165.190.79","mesh":{"netbird":{"ip":"100.122.165.149"}}}'::jsonb,1,now(),now())`,
|
||||
uuid.New())
|
||||
|
||||
addr, _, err := ResolveHost(ctx, pool, "host:vps", DefaultUser)
|
||||
if err != nil {
|
||||
t.Fatalf("ResolveHost: %v", err)
|
||||
}
|
||||
if addr != "82.165.190.79" {
|
||||
t.Errorf("addr = %q, want public_ipv4 (mesh unreachable from container)", addr)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResolveExecTargetForCheckLXCRoutesViaHost(t *testing.T) {
|
||||
// An LXC guest with a `hosts` edge to a proxmox host must resolve to the
|
||||
// HOST's address (the host-hop target), wrapped as `pct exec`.
|
||||
pool := newRemotePool(t)
|
||||
ctx := context.Background()
|
||||
hostID := uuid.New()
|
||||
guestID := uuid.New()
|
||||
mustExec(t, pool, ctx, `INSERT INTO entities (id, slug, type, name, state, attributes, version, created_at, updated_at)
|
||||
VALUES ($1,'host:hubris','proxmox-host','hubris','active','{"lan_ip":"192.168.8.77"}'::jsonb,1,now(),now())`, hostID)
|
||||
mustExec(t, pool, ctx, `INSERT INTO entities (id, slug, type, name, state, attributes, version, created_at, updated_at)
|
||||
VALUES ($1,'lxc:rclone','lxc','rclone','active','{"pve_id":"132"}'::jsonb,1,now(),now())`, guestID)
|
||||
mustExec(t, pool, ctx, `INSERT INTO relationships (source_id, target_id, type, valid_from, created_at)
|
||||
VALUES ($1,$2,'hosts',now(),now())`, hostID, guestID)
|
||||
|
||||
et, err := ResolveExecTargetForCheck(ctx, pool, guestID, "lxc", DefaultUser)
|
||||
if err != nil {
|
||||
t.Fatalf("ResolveExecTargetForCheck: %v", err)
|
||||
}
|
||||
if et.Host != "192.168.8.77" {
|
||||
t.Errorf("Host = %q, want proxmox host lan_ip 192.168.8.77 (host-hop)", et.Host)
|
||||
}
|
||||
out := et.Wrap("/opt/oikos/checks/cpu_check.sh")
|
||||
if !strings.Contains(out, "pct exec 132") {
|
||||
t.Errorf("guest wrap must use pct exec 132, got %q", out)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResolveExecTargetForCheckHostIsDirect(t *testing.T) {
|
||||
// A host-like target resolves to its own address with identity wrap.
|
||||
pool := newRemotePool(t)
|
||||
ctx := context.Background()
|
||||
hid := uuid.New()
|
||||
mustExec(t, pool, ctx, `INSERT INTO entities (id, slug, type, name, state, attributes, version, created_at, updated_at)
|
||||
VALUES ($1,'ws:mini','workstation','mini','active','{"lan_ip":"192.168.178.182","user":"dtoro"}'::jsonb,1,now(),now())`, hid)
|
||||
|
||||
et, err := ResolveExecTargetForCheck(ctx, pool, hid, "workstation", DefaultUser)
|
||||
if err != nil {
|
||||
t.Fatalf("ResolveExecTargetForCheck: %v", err)
|
||||
}
|
||||
if et.Host != "192.168.178.182" {
|
||||
t.Errorf("Host = %q, want 192.168.178.182", et.Host)
|
||||
}
|
||||
// Workstation's top-level `user` must be honored (the mac-mini fix).
|
||||
if et.User != "dtoro" {
|
||||
t.Errorf("User = %q, want dtoro (top-level user attr)", et.User)
|
||||
}
|
||||
if cmd := et.Wrap("uptime"); cmd != "uptime" {
|
||||
t.Errorf("host wrap must be identity, got %q", cmd)
|
||||
}
|
||||
}
|
||||
68
internal/remote/testutil_test.go
Normal file
68
internal/remote/testutil_test.go
Normal file
@@ -0,0 +1,68 @@
|
||||
package remote
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/dtoro/oikos/internal/db"
|
||||
"github.com/jackc/pgx/v5"
|
||||
)
|
||||
|
||||
// createTestDB provisions a throwaway migrated database off baseURL, the same
|
||||
// convention as internal/scheduler/coverage_test.go. The base URL must point
|
||||
// at a Postgres superuser-capable connection.
|
||||
func createTestDB(t *testing.T, baseURL string) *db.Pool {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
|
||||
admin, err := pgx.Connect(ctx, baseURL)
|
||||
if err != nil {
|
||||
t.Fatalf("connect admin: %v", err)
|
||||
}
|
||||
dbName := fmt.Sprintf("oikos_rem_%08x", rand.Int63())
|
||||
if _, err := admin.Exec(ctx, "CREATE DATABASE "+dbName); err != nil {
|
||||
admin.Close(ctx)
|
||||
t.Fatalf("create test db: %v", err)
|
||||
}
|
||||
admin.Close(ctx)
|
||||
|
||||
at := strings.LastIndex(baseURL, "/")
|
||||
testURL := baseURL[:at+1] + dbName
|
||||
if q := strings.Index(baseURL[at:], "?"); q >= 0 {
|
||||
testURL += baseURL[at+q:]
|
||||
}
|
||||
|
||||
pool, err := db.New(ctx, testURL)
|
||||
if err != nil {
|
||||
t.Fatalf("connect test db: %v", err)
|
||||
}
|
||||
if err := pool.Migrate(ctx); err != nil {
|
||||
t.Fatalf("migrate: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
pool.Close()
|
||||
if admin, err := pgx.Connect(ctx, baseURL); err == nil {
|
||||
admin.Exec(ctx, "DROP DATABASE IF EXISTS "+dbName+" WITH (FORCE)")
|
||||
admin.Close(ctx)
|
||||
}
|
||||
})
|
||||
return pool
|
||||
}
|
||||
|
||||
func getenvOrDefault(key, def string) string {
|
||||
if v := os.Getenv(key); v != "" {
|
||||
return v
|
||||
}
|
||||
return def
|
||||
}
|
||||
|
||||
func mustExec(t *testing.T, pool *db.Pool, ctx context.Context, q string, args ...any) {
|
||||
t.Helper()
|
||||
if _, err := pool.Exec(ctx, q, args...); err != nil {
|
||||
t.Fatalf("exec %s: %v", q, err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user