Compare commits

...
18 Commits
Author SHA1 Message Date
cabbage 620f882fe2 feat: isolate simultaneous RDP fixture profiles 2026-09-14 07:16:39 +00:00
cabbage 9a2b0fd3c1 docs: clarify RDP recovery and prerequisites 2026-09-14 07:08:41 +00:00
cabbage e3aadd8dd6 docs: make RDP fixture management self-service 2026-09-14 07:07:23 +00:00
cabbage a68729def0 fix: fail closed on unowned IPv6 listener 2026-09-14 06:53:49 +00:00
cabbage 3fdd0e2774 feat: add Guacamole stale-session repair action 2026-09-14 06:50:32 +00:00
cabbage cbd3f40523 fix: recognize owned Guacamole IPv6 forward 2026-09-14 06:46:21 +00:00
cabbage 102a5486f4 fix: make Guacamole fixture dual-stack and keep VM display awake 2026-09-14 06:31:07 +00:00
cabbage d01bb5a9a7 test: wait for Guacamole readiness 2026-09-14 04:34:18 +00:00
cabbage 38e6ca5eed fix: keep Guacamole VRDE tunnel alive 2026-09-14 04:02:32 +00:00
cabbage 6882a6808a feat: add interactive Windows installer 2026-09-11 10:10:48 +00:00
cabbage 9e181b81dc test: cover daemon telemetry publication 2026-09-11 09:43:24 +00:00
cabbage c5e9bc8b35 test: cover cursor decoding and script integrity metrics 2026-09-11 09:39:33 +00:00
cabbage 0e18af02bb feat: publish durable queue and spool gauges 2026-09-11 09:28:54 +00:00
cabbage ea8a4d8846 build: embed Windows release resources 2026-09-11 09:24:32 +00:00
cabbage 6c35c1a1d3 build: publish linux arm64 release artifacts 2026-09-11 09:18:29 +00:00
cabbage 15b2c4f9d3 fix: map release staging into toolchain workspace 2026-09-11 09:16:58 +00:00
cabbage cc31f6cb67 build: add reproducible release bundles 2026-09-11 09:13:36 +00:00
cabbage 2cf563d88e feat: add bounded operational metrics 2026-09-11 09:09:52 +00:00
40 changed files with 1950 additions and 79 deletions
+41
View File
@@ -26,12 +26,22 @@ import (
"google.golang.org/grpc"
)
// buildVersion is set by scripts/release for published artifacts. Development
// and test builds intentionally retain the explicit non-release value.
var buildVersion = "dev"
func main() {
var configPath string
var checkConfig bool
var version bool
flag.StringVar(&configPath, "config", "", "absolute server TOML configuration path")
flag.BoolVar(&checkConfig, "check-config", false, "validate server configuration and exit")
flag.BoolVar(&version, "version", false, "print build version and exit")
flag.Parse()
if version {
fmt.Fprintln(os.Stdout, buildVersion)
return
}
if configPath == "" {
log.Print("rvbox-server: --config is required")
os.Exit(2)
@@ -99,6 +109,7 @@ func run(configPath string) error {
}
recoveryContext, cancelRecovery := context.WithCancel(context.Background())
defer cancelRecovery()
go publishServerTelemetry(recoveryContext, persistence, health)
// Recovery is deliberately asynchronous: liveness and incident inspection
// remain available while committed-range checks run. Readiness becomes
// true only after the real SQLite/segment recovery completes.
@@ -134,6 +145,7 @@ func run(configPath string) error {
controlService, err := control.NewService(control.Options{
Store: persistence, DefaultQueueTTL: configured.Queue.DefaultTTL, DefaultQueueTTLSet: true, TakeoverTTL: configured.Protocol.TakeoverTTL,
WakeClient: registry.Wake,
Metrics: health,
Limits: agentproto.Limits{MaxEnvelopeBytes: configured.Protocol.MaxAgentEnvelopeBytes, MaxExecutionSpecBytes: configured.Protocol.MaxExecutionSpecBytes, MaxRawChunkBytes: configured.Protocol.MaxRawChunkBytes, MaxScriptBytes: configured.Protocol.MaxScriptBytes, MaxDetailBytes: configured.Storage.ProtocolDetailMaxBytes},
})
if err != nil {
@@ -156,6 +168,7 @@ func run(configPath string) error {
}
agent := &session.AgentServer{
Store: persistence, Registry: registry, Path: configured.Server.AgentPath,
Metrics: health,
Limits: agentproto.Limits{
MaxEnvelopeBytes: configured.Protocol.MaxAgentEnvelopeBytes,
MaxExecutionSpecBytes: configured.Protocol.MaxExecutionSpecBytes,
@@ -227,3 +240,31 @@ func run(configPath string) error {
return httpErr
}
}
func publishServerTelemetry(ctx context.Context, persistence *store.Store, health *observability.Health) {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for {
publishServerTelemetrySample(ctx, persistence.Telemetry, health)
select {
case <-ctx.Done():
return
case <-ticker.C:
}
}
}
// publishServerTelemetrySample performs one intentionally bounded, best-effort
// telemetry read. A failed metrics read must never affect server readiness or
// durable command processing.
func publishServerTelemetrySample(ctx context.Context, read func(context.Context) (store.Telemetry, error), health *observability.Health) {
sampleContext, cancel := context.WithTimeout(ctx, time.Second)
telemetry, err := read(sampleContext)
cancel()
if err != nil {
health.Inc("telemetry_read_failure")
return
}
health.SetGauge("queue_depth", float64(telemetry.QueuedCommands))
health.SetGauge("server_charged_bytes", float64(telemetry.ChargedBytes))
}
+34
View File
@@ -0,0 +1,34 @@
package main
import (
"context"
"errors"
"testing"
"github.com/rvbox/rvbox/internal/observability"
"github.com/rvbox/rvbox/internal/server/store"
)
func TestPublishServerTelemetrySample_HP_OPS_07(t *testing.T) {
t.Parallel()
health := observability.New()
publishServerTelemetrySample(context.Background(), func(context.Context) (store.Telemetry, error) {
return store.Telemetry{QueuedCommands: 3, ChargedBytes: 4096}, nil
}, health)
snapshot := health.MetricsSnapshot()
if snapshot.Gauges["queue_depth"] != 3 || snapshot.Gauges["server_charged_bytes"] != 4096 {
t.Fatalf("server telemetry gauges = %#v", snapshot.Gauges)
}
}
func TestPublishServerTelemetrySampleFailure_BH_OPS_04(t *testing.T) {
t.Parallel()
health := observability.New()
publishServerTelemetrySample(context.Background(), func(context.Context) (store.Telemetry, error) {
return store.Telemetry{}, errors.New("store unavailable")
}, health)
_, _, counters := health.Snapshot()
if counters["telemetry_read_failure"] != 1 {
t.Fatalf("telemetry failure counters = %#v", counters)
}
}
+52
View File
@@ -0,0 +1,52 @@
package main
import (
"fmt"
"strconv"
"strings"
"github.com/rvbox/rvbox/internal/config"
)
const (
defaultInstallerStateDir = `C:\ProgramData\RVBox\data`
// Public is a pre-existing directory usable by the active user, LocalSystem,
// and LocalService fallback. Durable RVBox state stays separately protected.
defaultInstallerWorkDir = `C:\Users\Public`
)
// packagedClientConfig creates the deliberately small first-install TOML. The
// strict production decoder validates it before an elevated installer writes
// anything, so the SCM service never starts from a hand-built invalid file.
func packagedClientConfig(serverURL, clientID string) ([]byte, error) {
serverURL = strings.TrimSpace(serverURL)
clientID = installerClientID(clientID)
if serverURL == "" {
return nil, fmt.Errorf("this RVBox package has no bootstrap server URL")
}
content := fmt.Sprintf("# Generated by the RVBox installer. Edit only through a deliberate reconfiguration.\n[client]\nserver_url = %s\nstate_dir = %s\nclient_id = %s\ndaemon_cwd = %s\n\n[observability]\n# Disabled by default; enable a loopback listener only when diagnostics are needed.\nlisten = \"\"\nlog_file = \"\"\n", strconv.Quote(serverURL), strconv.Quote(defaultInstallerStateDir), strconv.Quote(clientID), strconv.Quote(defaultInstallerWorkDir))
if _, err := config.DecodeClient([]byte(content), config.ClientOptions{Platform: config.PlatformWindows}); err != nil {
return nil, fmt.Errorf("validate packaged client configuration: %w", err)
}
return []byte(content), nil
}
func installerClientID(hostname string) string {
hostname = strings.TrimSpace(hostname)
var result strings.Builder
for _, value := range []byte(hostname) {
if value >= 0x21 && value <= 0x7e {
result.WriteByte(value)
} else {
result.WriteByte('-')
}
}
value := strings.Trim(result.String(), "-")
if value == "" {
value = "rvbox-client"
}
if len(value) > 128 {
value = value[:128]
}
return value
}
+16
View File
@@ -0,0 +1,16 @@
//go:build !windows
package main
import (
"errors"
"io"
)
func runInteractiveInstaller(io.Writer, io.Writer) error {
return errors.New("the v1 interactive installer is implemented for Windows only")
}
func installPackagedDefault(io.Writer) error {
return errors.New("the v1 packaged installer is implemented for Windows only")
}
+207
View File
@@ -0,0 +1,207 @@
//go:build windows
package main
import (
"errors"
"fmt"
"io"
"os"
"path/filepath"
"strings"
"syscall"
"unsafe"
"github.com/rvbox/rvbox/internal/client/windowsservice"
"golang.org/x/sys/windows"
)
const (
installerInfoIcon = 0x00000040
installerErrorIcon = 0x00000010
)
var (
installerUser32 = syscall.NewLazyDLL("user32.dll")
installerMessageBox = installerUser32.NewProc("MessageBoxW")
)
// runInteractiveInstaller is the double-click entry point. It does no
// machine-wide mutation itself: ShellExecute's runas verb causes Windows to
// show the standard UAC consent prompt before the elevated child installs.
func runInteractiveInstaller(_ io.Writer, _ io.Writer) error {
executable, err := os.Executable()
if err != nil {
showInstallerMessage("RVBox setup could not locate its executable.", installerErrorIcon)
return err
}
file, err := windows.UTF16PtrFromString(executable)
if err != nil {
return err
}
arguments, err := windows.UTF16PtrFromString("--install-default")
if err != nil {
return err
}
if err := windows.ShellExecute(0, installerUTF16("runas"), file, arguments, nil, windows.SW_SHOWNORMAL); err != nil {
showInstallerMessage("RVBox setup was cancelled or could not obtain administrator approval.\n\nNo changes were made.", installerErrorIcon)
return fmt.Errorf("request installer elevation: %w", err)
}
return nil
}
// installPackagedDefault executes only in the UAC-elevated child. It preserves
// a pre-existing operator configuration, atomically updates the executable,
// installs/updates the SCM service, and starts it.
func installPackagedDefault(_ io.Writer) error {
err := installPackagedDefaultFiles()
if err != nil {
showInstallerMessage("RVBox could not be installed.\n\n"+err.Error()+"\n\nNo existing configuration was replaced.", installerErrorIcon)
return err
}
showInstallerMessage("RVBox is installed and its Windows service is starting.\n\nThe notification-area icon starts automatically at the next sign-in.", installerInfoIcon)
return nil
}
func installPackagedDefaultFiles() error {
if strings.TrimSpace(bootstrapServerURL) == "" {
return errors.New("this RVBox package was built without a bootstrap server URL")
}
source, err := os.Executable()
if err != nil {
return fmt.Errorf("locate installer executable: %w", err)
}
programFiles := os.Getenv("ProgramFiles")
if programFiles == "" {
programFiles = `C:\Program Files`
}
programData := os.Getenv("ProgramData")
if programData == "" {
programData = `C:\ProgramData`
}
targetDirectory := filepath.Join(programFiles, "RVBox")
targetExecutable := filepath.Join(targetDirectory, "rvbox.exe")
configDirectory := filepath.Join(programData, "RVBox")
configPath := filepath.Join(configDirectory, "client.toml")
// A running service can keep the prior executable open. Stop it before the
// atomic replacement; a missing service is already a successful first run.
if err := windowsservice.Stop(30); err != nil {
return fmt.Errorf("stop existing RVBox service: %w", err)
}
if err := os.MkdirAll(targetDirectory, 0o755); err != nil {
return fmt.Errorf("create installation directory: %w", err)
}
if !sameWindowsPath(source, targetExecutable) {
if err := copyInstallerExecutable(source, targetExecutable); err != nil {
return err
}
}
if err := os.MkdirAll(configDirectory, 0o700); err != nil {
return fmt.Errorf("create configuration directory: %w", err)
}
if _, err := os.Stat(configPath); errors.Is(err, os.ErrNotExist) {
hostname, hostnameErr := os.Hostname()
if hostnameErr != nil {
hostname = "rvbox-client"
}
content, configErr := packagedClientConfig(bootstrapServerURL, hostname)
if configErr != nil {
return configErr
}
if err := writeMachineConfig(configPath, content); err != nil {
return err
}
} else if err != nil {
return fmt.Errorf("inspect existing configuration: %w", err)
}
if err := windowsservice.Install(windowsservice.InstallSpec{ExecutablePath: targetExecutable, ConfigPath: configPath, Startup: windowsservice.StartupAutomatic}); err != nil {
return fmt.Errorf("install and start RVBox service: %w", err)
}
return nil
}
func copyInstallerExecutable(source, target string) error {
input, err := os.Open(source)
if err != nil {
return fmt.Errorf("open installer executable: %w", err)
}
defer input.Close()
temporary := target + ".new"
_ = os.Remove(temporary)
output, err := os.OpenFile(temporary, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o755)
if err != nil {
return fmt.Errorf("create replacement executable: %w", err)
}
_, copyErr := io.Copy(output, input)
if syncErr := output.Sync(); copyErr == nil {
copyErr = syncErr
}
if closeErr := output.Close(); copyErr == nil {
copyErr = closeErr
}
if copyErr != nil {
_ = os.Remove(temporary)
return fmt.Errorf("copy installer executable: %w", copyErr)
}
from, fromErr := windows.UTF16PtrFromString(temporary)
to, toErr := windows.UTF16PtrFromString(target)
if fromErr != nil || toErr != nil {
_ = os.Remove(temporary)
return errors.New("encode replacement executable path")
}
if err := windows.MoveFileEx(from, to, windows.MOVEFILE_REPLACE_EXISTING|windows.MOVEFILE_WRITE_THROUGH); err != nil {
_ = os.Remove(temporary)
return fmt.Errorf("activate replacement executable: %w", err)
}
return nil
}
func writeMachineConfig(path string, content []byte) error {
temporary := path + ".new"
_ = os.Remove(temporary)
if err := os.WriteFile(temporary, content, 0o600); err != nil {
return fmt.Errorf("write default configuration: %w", err)
}
from, fromErr := windows.UTF16PtrFromString(temporary)
to, toErr := windows.UTF16PtrFromString(path)
if fromErr != nil || toErr != nil {
_ = os.Remove(temporary)
return errors.New("encode configuration path")
}
if err := windows.MoveFileEx(from, to, windows.MOVEFILE_REPLACE_EXISTING|windows.MOVEFILE_WRITE_THROUGH); err != nil {
_ = os.Remove(temporary)
return fmt.Errorf("activate default configuration: %w", err)
}
return protectMachineConfig(path)
}
func protectMachineConfig(path string) error {
descriptor, err := windows.SecurityDescriptorFromString("D:P(A;;FA;;;SY)(A;;FA;;;BA)")
if err != nil {
return err
}
dacl, _, err := descriptor.DACL()
if err != nil {
return err
}
if err := windows.SetNamedSecurityInfo(path, windows.SE_FILE_OBJECT, windows.DACL_SECURITY_INFORMATION|windows.PROTECTED_DACL_SECURITY_INFORMATION, nil, nil, dacl, nil); err != nil {
return fmt.Errorf("protect generated configuration: %w", err)
}
return nil
}
func sameWindowsPath(first, second string) bool {
return strings.EqualFold(filepath.Clean(first), filepath.Clean(second))
}
func installerUTF16(value string) *uint16 {
encoded, _ := windows.UTF16PtrFromString(value)
return encoded
}
func showInstallerMessage(message string, icon uintptr) {
caption := installerUTF16("RVBox Setup")
text := installerUTF16(message)
_, _, _ = installerMessageBox.Call(0, uintptr(unsafe.Pointer(text)), uintptr(unsafe.Pointer(caption)), icon)
}
+58 -10
View File
@@ -26,6 +26,9 @@ import (
"google.golang.org/protobuf/types/known/timestamppb"
)
// buildVersion is set by scripts/release for published artifacts.
var buildVersion = "dev"
func main() {
if err := run(os.Args[1:], os.Stdout, os.Stderr); err != nil {
fmt.Fprintln(os.Stderr, "rvbox:", err)
@@ -35,15 +38,24 @@ func main() {
var nativeTestContextFailures map[clientwindows.ExecutionContext]bool
// bootstrapServerURL is set only for a packaged Windows installer. A fresh
// machine has no service connection from which it could discover this value,
// so release packaging must deliberately supply the intended WSS endpoint.
var bootstrapServerURL string
// run is deliberately a small mode dispatcher. The SCM service invokes only
// --service with an explicit config path; tray/helper modes cannot silently
// turn an ordinary process invocation into a privileged service.
func run(args []string, output, diagnostics io.Writer) error {
if len(args) == 1 && args[0] == "--version" {
_, err := fmt.Fprintln(output, buildVersion)
return err
}
if len(args) == 0 {
return errors.New("an internal mode is required (use --help)")
return runInteractiveInstaller(output, diagnostics)
}
if args[0] == "--help" || args[0] == "-h" {
_, err := io.WriteString(output, "usage: rvbox --service|--tray|--launcher|--signal-helper|--check-config|--install-service|--uninstall-service|--configure-service|--start-service|--stop-service|--restart-service --config PATH\n")
_, err := io.WriteString(output, "usage: rvbox [--install-default]|--service|--tray|--launcher|--signal-helper|--check-config|--install-service|--uninstall-service|--configure-service|--start-service|--stop-service|--restart-service --config PATH\n")
return err
}
flags := flag.NewFlagSet("rvbox", flag.ContinueOnError)
@@ -56,6 +68,7 @@ func run(args []string, output, diagnostics io.Writer) error {
channel := flags.String("channel", "", "private launcher channel (internal use only)")
checkConfig := flags.Bool("check-config", false, "validate client configuration and exit")
install := flags.Bool("install-service", false, "install or update the machine-wide service")
installDefault := flags.Bool("install-default", false, "install the packaged default configuration (internal installer use)")
uninstall := flags.Bool("uninstall-service", false, "remove the machine-wide service")
configure := flags.Bool("configure-service", false, "configure machine-wide service startup mode")
startup := flags.String("startup", string(windowsservice.StartupAutomatic), "service startup mode: automatic or manual")
@@ -70,7 +83,7 @@ func run(args []string, output, diagnostics io.Writer) error {
return fmt.Errorf("unexpected argument %q", flags.Arg(0))
}
selected := 0
for _, value := range []bool{*serviceMode, *trayMode, *launcherMode, *signalHelperMode, *checkConfig, *install, *uninstall, *configure, *start, *stop, *restart} {
for _, value := range []bool{*serviceMode, *trayMode, *launcherMode, *signalHelperMode, *checkConfig, *install, *installDefault, *uninstall, *configure, *start, *stop, *restart} {
if value {
selected++
}
@@ -93,6 +106,9 @@ func run(args []string, output, diagnostics io.Writer) error {
}
return windowsservice.Install(windowsservice.InstallSpec{ExecutablePath: executable, ConfigPath: *configPath, Startup: windowsservice.StartupAutomatic})
}
if *installDefault {
return installPackagedDefault(diagnostics)
}
if *uninstall {
return windowsservice.Uninstall()
}
@@ -162,17 +178,20 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ
diagnostics = io.MultiWriter(formattedRuntimeLog, diagnostics)
}
health := observability.New()
go func() {
if serveErr := health.Serve(ctx, configured.Observability.Listen, observability.Paths{Liveness: configured.Observability.LivenessPath, Readiness: configured.Observability.ReadinessPath, Metrics: configured.Observability.MetricsPath}); serveErr != nil && ctx.Err() == nil && diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox client observability endpoint stopped: %v\n", serveErr)
}
}()
if configured.Observability.Listen != "" {
go func() {
if serveErr := health.Serve(ctx, configured.Observability.Listen, observability.Paths{Liveness: configured.Observability.LivenessPath, Readiness: configured.Observability.ReadinessPath, Metrics: configured.Observability.MetricsPath}); serveErr != nil && ctx.Err() == nil && diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox client observability endpoint stopped: %v\n", serveErr)
}
}()
}
state, err := spool.Open(ctx, spool.Options{DataDir: configured.Client.StateDir, BusyTimeout: 5 * time.Second, TombstoneLimit: configured.Storage.TombstoneMaxEntries, MaxScriptBytes: configured.Execution.MaxScriptBytes, MaxExecutionSpecBytes: configured.Execution.MaxExecutionSpecBytes, QuotaLimits: spool.QuotaLimits{HardAllocationBytes: 1 << 20, CommandOutputBytes: configured.Storage.CommandOutputLimitBytes, CommandTotalBytes: configured.Storage.CommandTotalLimitBytes, ClientTotalBytes: configured.Storage.ClientTotalLimitBytes, CloseoutReserveBytes: configured.Storage.CommandCloseoutReserveBytes}})
if err != nil {
health.SetDirty(true)
return fmt.Errorf("open client durable state: %w", err)
}
defer state.Close()
go publishClientTelemetry(ctx, state, health)
if diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox client service initialized for %s\n", configured.Client.ServerURL)
}
@@ -190,7 +209,7 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ
if err != nil {
return fmt.Errorf("configure command supervisor: %w", err)
}
executor, err := agent.NewExecutor(agent.ExecutorOptions{Store: state, Supervisor: supervised, WorkDir: configured.Client.DaemonCWD, Notify: func(issue domain.UUID) {
executor, err := agent.NewExecutor(agent.ExecutorOptions{Store: state, Supervisor: supervised, WorkDir: configured.Client.DaemonCWD, Metrics: health, Notify: func(issue domain.UUID) {
select {
case eventReady <- issue:
default:
@@ -216,7 +235,6 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ
} else if len(recovered) > 0 && diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox recovered %d uncertain launch(es)\n", len(recovered))
}
health.SetReady(true)
if runErr := agent.Run(ctx, agent.RunnerOptions{
Store: state,
Dial: func(dialContext context.Context) (agent.Transport, error) {
@@ -227,6 +245,8 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ
Jitter: agent.CryptoJitter, Now: func() time.Time { return time.Now().UTC() },
OnDispatch: executor.Dispatch, OnScriptReady: executor.ScriptReady, OnStdin: executor.Stdin,
OnCloseStdin: executor.CloseStdin, OnSignal: executor.Signal, OnTerminate: executor.Terminate, EventReady: eventReady,
Metrics: health,
OnSessionActive: health.SetReady,
OnSessionError: func(sessionErr error) {
if diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox client session retry: %v\n", sessionErr)
@@ -242,6 +262,34 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ
return nil
}
func publishClientTelemetry(ctx context.Context, state *spool.Store, health *observability.Health) {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for {
publishClientTelemetrySample(ctx, state.Telemetry, health)
select {
case <-ctx.Done():
return
case <-ticker.C:
}
}
}
// publishClientTelemetrySample performs one intentionally bounded, best-effort
// telemetry read. An unavailable spool metric must not interrupt reconnect or
// command supervision.
func publishClientTelemetrySample(ctx context.Context, read func(context.Context) (spool.Telemetry, error), health *observability.Health) {
sampleContext, cancel := context.WithTimeout(ctx, time.Second)
telemetry, err := read(sampleContext)
cancel()
if err != nil {
health.Inc("telemetry_read_failure")
return
}
health.SetGauge("queue_depth", float64(telemetry.QueuedCommands))
health.SetGauge("client_spool_bytes", float64(telemetry.ChargedBytes))
}
func clientJobProfiles(profiles config.Profiles) map[string]clientwindows.JobProfile {
return map[string]clientwindows.JobProfile{
rvboxv1.ExecutionProfile_EXECUTION_PROFILE_LIGHT.String(): toJobProfile(profiles.Light),
+70
View File
@@ -2,11 +2,16 @@ package main
import (
"bytes"
"context"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/rvbox/rvbox/internal/client/spool"
"github.com/rvbox/rvbox/internal/config"
"github.com/rvbox/rvbox/internal/observability"
)
func TestClientModeSelectionRequiresExactlyOneMode_HP_WINCLI_01(t *testing.T) {
@@ -23,6 +28,14 @@ func TestClientModeSelectionRequiresExactlyOneMode_HP_WINCLI_01(t *testing.T) {
}
}
func TestClientVersionDoesNotRequireConfiguration_HP_WINCLI_03(t *testing.T) {
t.Parallel()
var output, diagnostics bytes.Buffer
if err := run([]string{"--version"}, &output, &diagnostics); err != nil || output.String() != "dev\n" {
t.Fatalf("version = %q, %v", output.String(), err)
}
}
func TestClientHTTPClientAcceptsMatchingSelfSignedLeaf(t *testing.T) {
t.Parallel()
server := httptest.NewTLSServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) {
@@ -69,3 +82,60 @@ func TestNonWindowsServiceModesRemainExplicitlyUnsupported_BH_WINCLI_01(t *testi
t.Fatal("non-Windows tray mode unexpectedly available")
}
}
func TestPackagedClientConfigUsesMinimalSafeDefaults_HP_WINCLI_04(t *testing.T) {
t.Parallel()
content, err := packagedClientConfig("wss://controller.example.test/v1/agent", "desktop-01")
if err != nil {
t.Fatal(err)
}
text := string(content)
for _, want := range []string{
`server_url = "wss://controller.example.test/v1/agent"`,
`client_id = "desktop-01"`,
`state_dir = "C:\\ProgramData\\RVBox\\data"`,
`daemon_cwd = "C:\\Users\\Public"`,
`listen = ""`,
} {
if !strings.Contains(text, want) {
t.Fatalf("generated configuration missing %q:\n%s", want, text)
}
}
if _, err := packagedClientConfig("", "desktop-01"); err == nil {
t.Fatal("missing bootstrap URL was accepted")
}
}
func TestInstallerClientIDNormalizesHostnames_HP_WINCLI_05(t *testing.T) {
t.Parallel()
if got := installerClientID(" desktop host\n"); got != "desktop-host" {
t.Fatalf("normalized hostname = %q", got)
}
if got := installerClientID("\x00"); got != "rvbox-client" {
t.Fatalf("fallback hostname = %q", got)
}
}
func TestPublishClientTelemetrySample_HP_OPS_07(t *testing.T) {
t.Parallel()
health := observability.New()
publishClientTelemetrySample(context.Background(), func(context.Context) (spool.Telemetry, error) {
return spool.Telemetry{QueuedCommands: 2, ChargedBytes: 2048}, nil
}, health)
snapshot := health.MetricsSnapshot()
if snapshot.Gauges["queue_depth"] != 2 || snapshot.Gauges["client_spool_bytes"] != 2048 {
t.Fatalf("client telemetry gauges = %#v", snapshot.Gauges)
}
}
func TestPublishClientTelemetrySampleFailure_BH_OPS_04(t *testing.T) {
t.Parallel()
health := observability.New()
publishClientTelemetrySample(context.Background(), func(context.Context) (spool.Telemetry, error) {
return spool.Telemetry{}, errors.New("spool unavailable")
}, health)
_, _, counters := health.Snapshot()
if counters["telemetry_read_failure"] != 1 {
t.Fatalf("telemetry failure counters = %#v", counters)
}
}
+7
View File
@@ -25,6 +25,9 @@ import (
const defaultControlSocket = "/run/rvbox/server.sock"
// buildVersion is set by scripts/release for published artifacts.
var buildVersion = "dev"
func main() {
if err := run(os.Args[1:], os.Stdout, os.Stderr); err != nil {
fmt.Fprintln(os.Stderr, "rvc:", err)
@@ -33,6 +36,10 @@ func main() {
}
func run(args []string, output, diagnostics io.Writer) error {
if len(args) == 1 && args[0] == "--version" {
_, err := fmt.Fprintln(output, buildVersion)
return err
}
if len(args) == 0 {
return errors.New("a command is required (stat or run)")
}
+8
View File
@@ -50,6 +50,14 @@ func TestGlobalRequestIDIsInjectedOnlyForMutations_HP_CTL_12(t *testing.T) {
}
}
func TestVersionDoesNotRequireControlSocket_HP_CTL_16(t *testing.T) {
t.Parallel()
var output, diagnostics bytes.Buffer
if err := run([]string{"--version"}, &output, &diagnostics); err != nil || output.String() != "dev\n" {
t.Fatalf("version = %q, %v", output.String(), err)
}
}
func TestRenderCommandStatShowsExpiryRetentionAndWindowsIdentity_HP_CTL_13(t *testing.T) {
t.Parallel()
expiry := time.Date(2026, 9, 6, 12, 0, 0, 0, time.UTC)
+6
View File
@@ -8,6 +8,7 @@ ARG PROTOC_AMD64_SHA256=a45cda0989c17dd950db55f6fbe1e5814c50fda08e87aa422980ac1f
ARG PROTOC_ARM64_SHA256=36b518ac14d90351cc6598228ed2bbe5afe4e357b1af470b07e0ec1609875de2
ARG PROTOC_GEN_GO_VERSION=1.36.12
ARG PROTOC_GEN_GO_GRPC_VERSION=1.6.2
ARG GO_WINRES_VERSION=0.3.3
ARG TARGETARCH
ENV GOCACHE=/cache/go-build \
@@ -34,6 +35,11 @@ RUN apt-get update \
RUN GOBIN=/usr/local/bin go install "google.golang.org/protobuf/cmd/protoc-gen-go@v${PROTOC_GEN_GO_VERSION}" \
&& GOBIN=/usr/local/bin go install "google.golang.org/grpc/cmd/protoc-gen-go-grpc@v${PROTOC_GEN_GO_GRPC_VERSION}"
# Release packaging generates a temporary, architecture-suffixed Windows
# resource object with this pinned pure-Go tool. It is intentionally part of
# the toolchain image rather than a host prerequisite.
RUN GOBIN=/usr/local/bin go install "github.com/tc-hib/go-winres@v${GO_WINRES_VERSION}"
RUN mkdir -p /cache/go-build /cache/go/pkg/mod /cache/xdg /workspace \
&& chown -R 1001:1001 /cache /workspace
+23 -3
View File
@@ -888,10 +888,28 @@ invisible Guest Control session must never answer it.
When that bounded manual step is necessary, use the tracked Docker-only
[`test/rdp-access`](../test/rdp-access/README.md) helper. It starts a
self-signed HTTPS Guacamole gateway only after `test-host prepare` holds the
fixture lease; VRDE remains loopback-only on Helium and its SSH tunnel is bound
fixture lease; VRDE remains loopback-only on Helium and its reconnecting SSH
watchdog/tunnel is bound
only to the helper's private Docker gateway. Stop the helper before the normal
`test-host reset`. It is a recovery interface, not a product component or a
replacement for Guest Control/native test automation.
The helper's `repair` action may restart only guacd when a dropped browser
leaves a stale worker holding the fixture's single VRDE slot; it must not
restart or mutate the Windows VM automatically.
Because Windows can power off the virtual monitor while the VM remains
running, `test-host prepare` must also reapply the disposable display/sleep
keepalive and record `prepare-display-keepalive`. A zero-bpp VRDE framebuffer
otherwise leaves an otherwise authenticated Guacamole browser at “Waiting for
response”. If the public browser hostname has an AAAA record, `rdp-access up
--bind 0.0.0.0` must own a tracked IPv6-to-IPv4 forward and expose its
`ipv6_forward` status; this prevents a browser selecting IPv6 for the WebSocket
from silently taking a different, refused path.
The helper preserves the single-fixture `default` profile for compatibility;
simultaneous already-running VM gateways use a unique `RDP_ACCESS_PROFILE` (or
`--profile`) with separate runtime/project state and HTTP/tunnel ports, while
`RVBOX_TEST_VBOX_HOST`, the documented VM/snapshot identity variables, and
`RDP_ACCESS_VRDE_PORT` select the actual target. Profile isolation does not
override the native controller's exclusive lease for prepare/reset operations.
For this provisioned lane, the approved host-only credential-file location is
`/home/cabbage/.local/share/rvbox-secrets/rvbox-win10-test.password`. It must
@@ -1157,8 +1175,10 @@ Use one declared test matrix and the same harness entry points everywhere:
native Windows integration suites, plus `smoke`, `idempotency`, `reconnect`,
and `retention` E2E scenarios;
- release candidates run every E2E scenario, including the resettable
interactive Windows VM matrix, Server Core smoke, destructive fault cases,
upgrade/recovery, and soak.
interactive Windows 10 VM matrix, destructive fault cases, upgrade/recovery,
and soak. Server Core, older-build, and ambiguous-multi-session smoke lanes
run when their dedicated fixtures are provisioned; they are explicitly
deferred compatibility work and do not block the Windows 10 v1 baseline.
CI allocates a run ID per job, always invokes `collect` after failure, and invokes
`reset` in an unconditional finalizer. Upload only the bounded redacted report,
+49
View File
@@ -25,6 +25,24 @@ the process is up; `/readyz` becomes successful only after durable recovery.
The public endpoint accepts only `wss://HOST/v1/agent`. Do not publish port
6900 or add a proxy route for JSON-RPC.
## Health, metrics, and logs
Keep the configured observability listener private to the host or monitoring
network. `/metrics` exposes only bounded, aggregate Prometheus samples: health
state; registration/takeover and protocol failures; session/reconnect and
heartbeat timing; dispatch/event/transition counts and latency; and client
output accounting. It never includes a command body, stdin, output, client ID,
or UUID label. Duration histograms use fixed buckets, so monitoring traffic
cannot create unbounded series.
For a Windows client, readiness is false while its durable spool is recovering
or it has no reconciled server session; it becomes true only while the active
WSS session can exchange command data. Server liveness starts before
asynchronous recovery, while server readiness remains false until that recovery
completes. The configured rotating JSON/text service logs are diagnostics, not
a command-output store; use `rvc` history and the audited storage for command
evidence.
## Backup and restore
Stop dispatch before copying data: stop the server gracefully, confirm it is
@@ -51,6 +69,37 @@ false and records an incident; do not delete segments to force readiness.
and full data directory, then start the known-good version. Preserve logs
and the failed copy for diagnosis.
## Release bundle
Build a release candidate only from a clean, committed worktree. The
containerized release wrapper embeds the supplied version in all three binaries,
creates an immutable `dist/rvbox-VERSION` directory, and writes `SHA256SUMS`
plus `manifest.json` only after the Windows executable has optionally been
signed:
```sh
scripts/release build \
--version 1.0.0-rc.1 \
--bootstrap-server-url wss://rvbox.example.test/v1/agent
(cd dist/rvbox-1.0.0-rc.1 && sha256sum -c SHA256SUMS)
dist/rvbox-1.0.0-rc.1/rvbox-server-linux-ARCH --version
dist/rvbox-1.0.0-rc.1/rvc-linux-ARCH --version
```
Replace `ARCH` with `amd64` or `arm64` for the Linux host. Both variants are
created from the same pinned toolchain and have the same embedded version.
For a Windows-signed release, provide an executable host-side signing hook via
`--sign-windows-hook /absolute/path/to/hook`. The wrapper invokes it with the
Windows executable path and version, then records `windows_signed: true` in the
manifest. Without that hook the manifest deliberately declares the artifact
unsigned; this is suitable for CI/test evidence but not a signed public
release. `--output` is intentionally constrained beneath this repository's
ignored `dist/` tree so the pinned build container always sees the exact output
mount. The wrapper never overwrites a final bundle, so correcting a failed
candidate requires choosing a new version/output or deliberately removing that
exact ignored `dist/` directory after preserving any evidence.
## Common incidents
- **No client / stale session:** verify nginx has WebSocket `101` entries for
+16 -4
View File
@@ -88,8 +88,16 @@ and collection. Do not expose the VM's RDP endpoints beyond the test LAN.
For the rare interactive UAC/manual-recovery step, use the Docker-only helper
in [`test/rdp-access`](../test/rdp-access/README.md). It creates a temporary
self-signed HTTPS Guacamole gateway while retaining VRDE on Helium loopback and
the SSH tunnel on a private Docker gateway. Follow its full lease/prepare/up/
the reconnecting SSH watchdog/tunnel on a private Docker gateway. Follow its
full lease/prepare/up/
down/reset lifecycle; it is not an alternative to the native test controller.
In public-bind mode the helper also owns a tracked IPv6-to-IPv4 `socat` forward
when the browser hostname has an AAAA record; check `ipv6_forward=active` in
`rdp-access status` before diagnosing a browser-side “Waiting for response”.
When more than one already-running VM needs browser access, set matching
`RVBOX_TEST_VBOX_*`/`RDP_ACCESS_VRDE_PORT` values and a unique
`RDP_ACCESS_PROFILE` plus HTTP/tunnel ports; the helper README has the
copy/paste example and explains the native-host lease boundary.
## Snapshots and reset contract
@@ -111,10 +119,14 @@ stateful run must:
1. Acquire the run lease and verify the VM name, UUID, and snapshot UUID.
2. Restore `baseline-clean-administrator` if the current state is not the baseline.
3. Start headless and wait for `VMState=running` plus Guest Additions readiness.
4. Run the bounded test, collect redacted artifacts, and close every Guest
4. Apply the disposable display/sleep keepalive (`monitor-timeout`, standby,
and hibernate timers set to zero) so VirtualBox VRDE cannot expose a
zero-bpp framebuffer after Windows idle timeout. This is recorded as
`prepare-display-keepalive` and is reapplied after every snapshot restore.
5. Run the bounded test, collect redacted artifacts, and close every Guest
Control process that was opened by the run.
5. Request a graceful guest shutdown and wait for `VMState=poweroff`.
6. Restore `baseline-clean-administrator` again and leave the VM powered off.
6. Request a graceful guest shutdown and wait for `VMState=poweroff`.
7. Restore `baseline-clean-administrator` again and leave the VM powered off.
Use `controlvm ... poweroff` only for a hung, disposable test; it can lose
guest state. Never delete any clean snapshot, unregister the VM, or alter
+10 -3
View File
@@ -208,6 +208,11 @@ Interactive browser access is a deliberately temporary recovery path only. See
[`test/rdp-access`](../test/rdp-access/README.md) for the Docker-only,
self-signed HTTPS Guacamole lifecycle; it must be started only after the native
fixture controller has prepared and leased the VM, and stopped before reset.
For a public hostname with an AAAA record, its `up --bind 0.0.0.0` mode also
tracks an IPv6-to-IPv4 forward; verify `ipv6_forward=active` before debugging
a browser stuck at “Waiting for response”. If several already-running VMs need
browser access at once, use the documented `RDP_ACCESS_PROFILE` and matching
host/VRDE/HTTP/tunnel-port overrides in the helper README.
The Linux production Compose asset has a separate, loopback-only smoke lane. It
uses a disposable self-signed key only under the ignored `.test-runs` tree,
@@ -286,6 +291,8 @@ execution context. The normal `rvboxtest` console session remains the subject
of active-user and elevation tests. See [testing-vm.md](testing-vm.md) for the
exact fixture contract.
The VM is the minimum smoke lane, so native multi-session/ambiguous-session,
Server Core, and older-build entries remain explicitly blocked until their own
fixtures exist. Their pure selector tests remain mandatory.
The VM is the current Windows 10 v1 smoke lane. Native
multi-session/ambiguous-session, Server Core, and older-build entries are
explicitly deferred compatibility work until their own fixtures exist; they do
not block the Windows 10 v1 baseline. Their pure selector tests remain
mandatory.
+24 -1
View File
@@ -12,6 +12,7 @@ import (
"github.com/rvbox/rvbox/internal/client/spool"
"github.com/rvbox/rvbox/internal/client/supervisor"
"github.com/rvbox/rvbox/internal/domain"
"github.com/rvbox/rvbox/internal/observability"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/timestamppb"
)
@@ -25,6 +26,7 @@ type Executor struct {
WorkDir string
Now func() time.Time
Notify func(domain.UUID)
Metrics *observability.Health
mu sync.Mutex
active map[domain.UUID]supervisor.Process
@@ -38,6 +40,7 @@ type ExecutorOptions struct {
WorkDir string
Now func() time.Time
Notify func(domain.UUID)
Metrics *observability.Health
}
func NewExecutor(options ExecutorOptions) (*Executor, error) {
@@ -47,7 +50,7 @@ func NewExecutor(options ExecutorOptions) (*Executor, error) {
if options.Now == nil {
options.Now = func() time.Time { return time.Now().UTC() }
}
return &Executor{Store: options.Store, Supervisor: options.Supervisor, WorkDir: options.WorkDir, Now: options.Now, Notify: options.Notify, active: make(map[domain.UUID]supervisor.Process), cancel: make(map[domain.UUID]context.CancelFunc), forced: make(map[domain.UUID]bool)}, nil
return &Executor{Store: options.Store, Supervisor: options.Supervisor, WorkDir: options.WorkDir, Now: options.Now, Notify: options.Notify, Metrics: options.Metrics, active: make(map[domain.UUID]supervisor.Process), cancel: make(map[domain.UUID]context.CancelFunc), forced: make(map[domain.UUID]bool)}, nil
}
// Dispatch is safe to invoke after CommandAccepted has been sent. Script
@@ -158,6 +161,7 @@ func (executor *Executor) launch(ctx context.Context, issue domain.UUID, revisio
executor.remove(issue)
return err
}
executor.incMetric("command_running")
executor.notify(issue)
go executor.watch(runContext, issue, revision, process, cancel)
return nil
@@ -171,6 +175,7 @@ func (executor *Executor) watch(ctx context.Context, issue domain.UUID, revision
break
}
if err != nil {
executor.incMetric("output_loss")
_ = executor.appendIncomplete(context.Background(), issue, revision, err.Error())
break
}
@@ -178,10 +183,13 @@ func (executor *Executor) watch(ctx context.Context, issue domain.UUID, revision
continue
}
if _, err := executor.Store.AppendOutput(context.Background(), issue, spool.OutputInput{Stream: chunk.Stream, Raw: chunk.Data, ObservedAt: executor.Now()}); err != nil {
executor.incMetric("output_loss")
_ = executor.appendIncomplete(context.Background(), issue, revision, err.Error())
_, _ = executor.Supervisor.Signal(context.Background(), process, supervisor.SignalKill)
break
}
executor.incMetric("output_chunk")
executor.incMetricBy("output_bytes", uint64(len(chunk.Data)))
executor.notify(issue)
}
status, waitErr := process.Wait(context.Background())
@@ -210,9 +218,24 @@ func (executor *Executor) appendLifecycle(ctx context.Context, issue domain.UUID
func (executor *Executor) appendLifecycleWithIdentity(ctx context.Context, issue domain.UUID, revision uint64, phase rvboxv1.CommandLifecycle, detail string, identity *rvboxv1.WindowsExecutionIdentity) error {
_, err := executor.Store.AppendLifecycleWithIdentity(ctx, issue, uint32(phase), revision, detail, executor.Now(), identity)
if err == nil {
executor.incMetric("command_transition")
}
return err
}
func (executor *Executor) incMetric(name string) {
if executor != nil && executor.Metrics != nil {
executor.Metrics.Inc(name)
}
}
func (executor *Executor) incMetricBy(name string, delta uint64) {
if executor != nil && executor.Metrics != nil {
executor.Metrics.IncBy(name, delta)
}
}
func (executor *Executor) reject(ctx context.Context, issue domain.UUID, revision uint64, cause error) error {
return executor.rejectWithIdentity(ctx, issue, revision, cause, nil)
}
+45
View File
@@ -12,6 +12,7 @@ import (
"github.com/rvbox/rvbox/internal/agentproto"
"github.com/rvbox/rvbox/internal/client/spool"
"github.com/rvbox/rvbox/internal/domain"
"github.com/rvbox/rvbox/internal/observability"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/timestamppb"
)
@@ -39,6 +40,12 @@ type RunnerOptions struct {
// transport session before normal reconnect backoff. It must not block; the
// durable spool and retry policy remain owned by Run.
OnSessionError func(error)
// OnSessionActive reports whether a reconciled transport is currently able
// to exchange command data. It lets daemon readiness represent connection
// state without making reconnect behavior depend on the observer.
OnSessionActive func(bool)
// Metrics is optional. It only receives static, bounded metric names.
Metrics *observability.Health
// EventReady wakes the active session after a supervisor worker appends a
// durable event. The network loop remains the sole writer; a reconnect can
// safely ignore a stale notification because replay reads the spool again.
@@ -98,6 +105,10 @@ func Run(ctx context.Context, options RunnerOptions) error {
if err := waitUntil(ctx, delay); err != nil {
return nil
}
if failures > 0 {
options.incMetric("reconnect_attempt")
}
options.incMetric("connection_attempt")
sessionStarted := options.Now()
sessionErr := runOnce(ctx, options)
if ctx.Err() != nil {
@@ -106,6 +117,9 @@ func Run(ctx context.Context, options RunnerOptions) error {
if sessionErr != nil && options.OnSessionError != nil {
options.OnSessionError(sessionErr)
}
if sessionErr != nil {
options.incMetric("connection_failure")
}
if options.Now().Sub(sessionStarted) >= options.Backoff.StableReset {
failures = 0
}
@@ -173,6 +187,7 @@ func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
}
session, err := Handshake(ctx, transport, options.Hello, limits)
if err != nil {
options.incMetric("protocol_error")
return fmt.Errorf("agent handshake: %w", err)
}
snapshot, err := options.Store.ReconcileSnapshot(ctx)
@@ -201,6 +216,17 @@ func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
if err := sendCapacity(ctx, transport, session, options.Hello, snapshot, limits); err != nil {
return fmt.Errorf("advertise client capacity: %w", err)
}
activeStarted := options.Now()
options.incMetric("session_established")
if options.OnSessionActive != nil {
options.OnSessionActive(true)
}
defer func() {
if options.OnSessionActive != nil {
options.OnSessionActive(false)
}
options.observeMetric("session_duration", options.Now().Sub(activeStarted))
}()
return serveActive(ctx, transport, options, session, sentEvents)
}
@@ -343,9 +369,11 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
}
envelope, err := agentproto.DecodeEnvelope(encoded, limits, rvboxv1.Platform_PLATFORM_WINDOWS)
if err != nil {
options.incMetric("protocol_error")
return fmt.Errorf("decode active agent frame: %w", err)
}
if envelope.GetSessionId() != session.ID || envelope.GetSessionGeneration() != session.Generation {
options.incMetric("stale_message")
return ErrProtocolHandshake
}
switch {
@@ -452,6 +480,11 @@ func flushEvents(ctx context.Context, transport Transport, store *spool.Store, s
func handleDispatch(ctx context.Context, transport Transport, options RunnerOptions, session Session, dispatch *rvboxv1.CommandDispatch, limits agentproto.Limits) error {
acceptance, err := PersistDispatch(ctx, options.Store, session, dispatch, options.Now(), limits)
accepted := err == nil
if accepted {
options.incMetric("command_dispatch_accepted")
} else {
options.incMetric("command_dispatch_rejected")
}
ack := &rvboxv1.CommandAccepted{IssueUuid: dispatch.GetIssueUuid(), CommandRevision: dispatch.GetCommandRevision(), Accepted: accepted}
if err != nil {
ack.Rejection = rejectionForError(err, dispatch.GetIssueUuid())
@@ -474,6 +507,18 @@ func handleDispatch(ctx context.Context, transport Transport, options RunnerOpti
return nil
}
func (options RunnerOptions) incMetric(name string) {
if options.Metrics != nil {
options.Metrics.Inc(name)
}
}
func (options RunnerOptions) observeMetric(name string, duration time.Duration) {
if options.Metrics != nil {
options.Metrics.ObserveDuration(name, duration)
}
}
func rejectionForError(err error, issue string) *rvboxv1.ControlError {
code := rvboxv1.ControlError_TRANSIENT
if errors.Is(err, spool.ErrCommandConflict) || errors.Is(err, spool.ErrAlreadyExecuted) {
+11 -1
View File
@@ -12,6 +12,7 @@ import (
"github.com/rvbox/rvbox/internal/agentproto"
"github.com/rvbox/rvbox/internal/client/spool"
"github.com/rvbox/rvbox/internal/domain"
"github.com/rvbox/rvbox/internal/observability"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/timestamppb"
)
@@ -46,7 +47,9 @@ func TestRunOnceReconcilesReplaysAndAdvertisesCapacity_HP_RUNTIME_01(t *testing.
reconcileResult, _ := proto.Marshal(&rvboxv1.AgentEnvelope{SessionId: "session", SessionGeneration: 1, Payload: &rvboxv1.AgentEnvelope_ReconcileResult{ReconcileResult: &rvboxv1.ReconcileResult{}}})
transport := &runnerTransport{reads: [][]byte{welcome, reconcileRequest, reconcileResult}, terminal: errors.New("transport closed")}
hello := &rvboxv1.ClientHello{ClientId: "runner-client", SupportedProtocol: &rvboxv1.ProtocolRange{Major: 1, MinMinor: 0, MaxMinor: 0}, DaemonVersion: "test", Platform: rvboxv1.Platform_PLATFORM_WINDOWS, Architecture: "amd64", DaemonCwd: `C:\`, SupportedShells: []rvboxv1.ShellType{rvboxv1.ShellType_SHELL_POWERSHELL}, ClientInstanceId: store.ClientInstanceID().String(), MaxRunningCommands: 1, MaxQueuedCommands: 1, SentAt: timestamppb.Now()}
err = RunOnce(ctx, RunnerOptions{Store: store, Dial: func(context.Context) (Transport, error) { return transport, nil }, Hello: hello, Limits: agentproto.DefaultLimits(), Backoff: BackoffOptions{Initial: time.Millisecond, Maximum: time.Millisecond, StableReset: time.Second}, Jitter: func(value time.Duration) time.Duration { return 0 }, Now: func() time.Time { return time.Now().UTC() }})
metrics := observability.New()
var active []bool
err = RunOnce(ctx, RunnerOptions{Store: store, Dial: func(context.Context) (Transport, error) { return transport, nil }, Hello: hello, Limits: agentproto.DefaultLimits(), Backoff: BackoffOptions{Initial: time.Millisecond, Maximum: time.Millisecond, StableReset: time.Second}, Jitter: func(value time.Duration) time.Duration { return 0 }, Now: func() time.Time { return time.Now().UTC() }, Metrics: metrics, OnSessionActive: func(value bool) { active = append(active, value) }})
if !errors.Is(err, transport.terminal) {
t.Fatalf("RunOnce error = %v, want transport close", err)
}
@@ -64,6 +67,13 @@ func TestRunOnceReconcilesReplaysAndAdvertisesCapacity_HP_RUNTIME_01(t *testing.
if err != nil || eventEnvelope.GetCommandEvent() == nil || eventEnvelope.GetCommandEvent().GetEventSeq() != 1 || eventEnvelope.GetCommandEvent().GetIssueUuid() != issue.String() {
t.Fatalf("replayed event = %#v, %v", eventEnvelope, err)
}
_, _, counters := metrics.Snapshot()
if counters["session_established"] != 1 {
t.Fatalf("session metrics = %#v", counters)
}
if len(active) != 2 || !active[0] || active[1] {
t.Fatalf("active transitions = %#v", active)
}
}
func TestRunnerOptionsRejectMissingJitter_HP_RUNTIME_02(t *testing.T) {
+34
View File
@@ -0,0 +1,34 @@
package spool
import (
"context"
"errors"
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
)
// Telemetry is the bounded durable-spool state exposed as client gauges. It
// contains no command body, output, UUID, or client identity.
type Telemetry struct {
QueuedCommands uint64
ChargedBytes uint64
}
func (store *Store) Telemetry(ctx context.Context) (Telemetry, error) {
if store == nil {
return Telemetry{}, errors.New("client spool is unavailable")
}
store.mu.Lock()
defer store.mu.Unlock()
if store.db == nil {
return Telemetry{}, errors.New("client spool is closed")
}
var result Telemetry
if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM commands WHERE terminal = 0 AND phase IN (?, ?)`, uint32(rvboxv1.CommandLifecycle_COMMAND_QUEUED), uint32(rvboxv1.CommandLifecycle_COMMAND_DISPATCHED)).Scan(&result.QueuedCommands); err != nil {
return Telemetry{}, err
}
if err := store.db.QueryRowContext(ctx, `SELECT client_total_charged_bytes FROM spool_counters WHERE singleton = 1`).Scan(&result.ChargedBytes); err != nil {
return Telemetry{}, err
}
return result, nil
}
+22
View File
@@ -0,0 +1,22 @@
package spool
import (
"context"
"path/filepath"
"testing"
"time"
)
func TestTelemetryReportsQueuedDurableSpool_HP_OPS_07(t *testing.T) {
t.Parallel()
ctx := context.Background()
store := openTestStore(t, ctx, filepath.Join(t.TempDir(), "spool"), DefaultTombstoneLimit)
issue := testUUID(t, "019c46f1-1d02-7000-8000-0000000000d2")
if _, err := store.AcceptCommand(ctx, testCommand(issue, []byte("telemetry")), time.Now().UTC()); err != nil {
t.Fatal(err)
}
telemetry, err := store.Telemetry(ctx)
if err != nil || telemetry.QueuedCommands != 1 || telemetry.ChargedBytes == 0 {
t.Fatalf("telemetry = %#v, %v", telemetry, err)
}
}
+14
View File
@@ -60,6 +60,20 @@ func TestClientDefaultsDecodeForWindowsWithoutNativeFilesystem(t *testing.T) {
}
}
func TestClientMayDisableObservabilityListener_HP_CFG_10(t *testing.T) {
t.Parallel()
client, err := DecodeClient([]byte("[client]\nstate_dir = \"C:\\\\ProgramData\\\\RVBox\\\\state\"\ndaemon_cwd = \"C:\\\\Users\\\\Public\"\n\n[observability]\nlisten = \"\"\nlog_file = \"\"\n"), ClientOptions{Platform: PlatformWindows})
if err != nil {
t.Fatalf("DecodeClient disabled observability: %v", err)
}
if client.Observability.Listen != "" || client.Observability.LogFile != "" {
t.Fatalf("disabled observability = %+v", client.Observability)
}
if _, err := DecodeServer([]byte("[observability]\nlisten = \"\"\n")); err == nil {
t.Fatal("server accepted a disabled observability listener")
}
}
func TestStrictTOMLRejections_HP_CFG_01(t *testing.T) {
t.Parallel()
+10 -5
View File
@@ -119,7 +119,7 @@ func validateServer(raw serverFile) (*Server, error) {
if err := protocolLimits(raw.Protocol.MaxAgentEnvelopeBytes, raw.Protocol.MaxExecutionSpecBytes, raw.Protocol.MaxRawChunkBytes, raw.Protocol.MaxScriptBytes, raw.Protocol.MaxControlRequestBytes, raw.Protocol.MaxJSONRPCBodyBytes); err != nil {
return nil, err
}
observability, err := validateObservability(raw.Observability, PlatformUnix)
observability, err := validateObservability(raw.Observability, PlatformUnix, false)
if err != nil {
return nil, err
}
@@ -220,7 +220,7 @@ func validateClient(raw clientFile, options ClientOptions) (*Client, error) {
if err != nil {
return nil, err
}
observability, err := validateObservability(raw.Observability, options.Platform)
observability, err := validateObservability(raw.Observability, options.Platform, true)
if err != nil {
return nil, err
}
@@ -439,9 +439,14 @@ func validateDeviceRate(key string, value uint64) error {
return nil
}
func validateObservability(raw observabilityFile, platform Platform) (Observability, error) {
if err := validateListener("observability.listen", raw.Listen); err != nil {
return Observability{}, err
func validateObservability(raw observabilityFile, platform Platform, allowDisabled bool) (Observability, error) {
if raw.Listen == "" && !allowDisabled {
return Observability{}, fmt.Errorf("observability.listen must be a host:port listener")
}
if raw.Listen != "" {
if err := validateListener("observability.listen", raw.Listen); err != nil {
return Observability{}, err
}
}
for name, value := range map[string]string{"liveness_path": raw.LivenessPath, "readiness_path": raw.ReadinessPath, "metrics_path": raw.MetricsPath} {
if err := validateHTTPPath("observability."+name, value); err != nil {
+21
View File
@@ -72,3 +72,24 @@ func TestCursorInputBounds_BH_CTL_01(t *testing.T) {
}
}
}
// FuzzDecodeCursorBounded_SEC_CTL_02 exercises the token boundary shared by
// every cursor-resumable control read. Decode must reject malformed or forged
// input without panicking or allocating past its documented token ceiling.
func FuzzDecodeCursorBounded_SEC_CTL_02(f *testing.F) {
codec, err := NewCursorCodec(bytes.Repeat([]byte{0x42}, 32))
if err != nil {
f.Fatal(err)
}
filters := HashCursorFilters([]byte("client=host-1\x00streams=stdout"))
valid, err := codec.Encode(Cursor{Kind: CursorKindOutput, FilterHash: filters, Position: []byte("event/offset"), SnapshotBoundary: []byte("upper-event")})
if err != nil {
f.Fatal(err)
}
for _, seed := range []string{"", "***", valid, strings.Repeat("a", 4096)} {
f.Add(seed)
}
f.Fuzz(func(t *testing.T, token string) {
_, _ = codec.Decode(token, CursorKindOutput, filters)
})
}
+134 -9
View File
@@ -6,6 +6,7 @@ package observability
import (
"context"
"fmt"
"math"
"net"
"net/http"
"sort"
@@ -27,9 +28,28 @@ type Health struct {
dirty atomic.Bool
mu sync.Mutex
count map[string]uint64
gauge map[string]float64
hist map[string]*durationHistogram
}
func New() *Health { return &Health{count: make(map[string]uint64)} }
type durationHistogram struct {
count uint64
sum float64
buckets [len(durationBuckets)]uint64
}
// The fixed buckets are intentionally shared by every duration metric. This
// keeps the Prometheus surface bounded and makes comparable operational
// latencies available without allowing callers to create arbitrary labels.
var durationBuckets = [...]float64{0.001, 0.005, 0.01, 0.05, 0.1, 0.25, 0.5, 1, 5, 15, 60}
func New() *Health {
return &Health{
count: make(map[string]uint64),
gauge: make(map[string]float64),
hist: make(map[string]*durationHistogram),
}
}
func (health *Health) SetReady(value bool) {
if health != nil {
@@ -44,11 +64,52 @@ func (health *Health) SetDirty(value bool) {
}
func (health *Health) Inc(name string) {
if health == nil || !validMetricName(name) {
health.IncBy(name, 1)
}
// IncBy records a non-negative integral counter delta. Metric names are
// deliberately restricted to a small identifier grammar; callers cannot turn
// a client ID, request UUID, or other untrusted data into a metric name.
func (health *Health) IncBy(name string, delta uint64) {
if health == nil || delta == 0 || !validMetricName(name) {
return
}
health.mu.Lock()
health.count[name]++
health.count[name] += delta
health.mu.Unlock()
}
// SetGauge records a bounded-cardinality instantaneous value. NaN and
// infinities are discarded because they are not safe Prometheus samples.
func (health *Health) SetGauge(name string, value float64) {
if health == nil || !validMetricName(name) || math.IsNaN(value) || math.IsInf(value, 0) {
return
}
health.mu.Lock()
health.gauge[name] = value
health.mu.Unlock()
}
// ObserveDuration records a duration in a fixed-bucket histogram. Negative
// durations are invalid (normally a caller clock error) and are ignored.
func (health *Health) ObserveDuration(name string, duration time.Duration) {
if health == nil || !validMetricName(name) || duration < 0 {
return
}
seconds := duration.Seconds()
health.mu.Lock()
histogram := health.hist[name]
if histogram == nil {
histogram = &durationHistogram{}
health.hist[name] = histogram
}
histogram.count++
histogram.sum += seconds
for index, upperBound := range durationBuckets {
if seconds <= upperBound {
histogram.buckets[index]++
}
}
health.mu.Unlock()
}
@@ -65,6 +126,46 @@ func (health *Health) Snapshot() (ready, dirty bool, counters map[string]uint64)
return health.ready.Load(), health.dirty.Load(), copyCounters
}
type MetricsSnapshot struct {
Ready bool
Dirty bool
Counters map[string]uint64
Gauges map[string]float64
Histograms map[string]DurationHistogramSnapshot
}
type DurationHistogramSnapshot struct {
Count uint64
Sum float64
Buckets [len(durationBuckets)]uint64
}
// MetricsSnapshot returns a copy suitable for rendering or testing. It never
// includes application identifiers or payloads.
func (health *Health) MetricsSnapshot() MetricsSnapshot {
if health == nil {
return MetricsSnapshot{Dirty: true}
}
health.mu.Lock()
defer health.mu.Unlock()
result := MetricsSnapshot{
Ready: health.ready.Load(), Dirty: health.dirty.Load(),
Counters: make(map[string]uint64, len(health.count)),
Gauges: make(map[string]float64, len(health.gauge)),
Histograms: make(map[string]DurationHistogramSnapshot, len(health.hist)),
}
for name, value := range health.count {
result.Counters[name] = value
}
for name, value := range health.gauge {
result.Gauges[name] = value
}
for name, value := range health.hist {
result.Histograms[name] = DurationHistogramSnapshot{Count: value.count, Sum: value.sum, Buckets: value.buckets}
}
return result
}
func (health *Health) Handler(paths Paths) http.Handler {
if paths.Liveness == "" {
paths.Liveness = "/livez"
@@ -92,17 +193,38 @@ func (health *Health) Handler(paths Paths) http.Handler {
_, _ = fmt.Fprintf(response, "ready=%s dirty=%s\n", strconv.FormatBool(ready), strconv.FormatBool(dirty))
})
mux.HandleFunc(paths.Metrics, func(response http.ResponseWriter, _ *http.Request) {
ready, dirty, counters := health.Snapshot()
snapshot := health.MetricsSnapshot()
response.Header().Set("Content-Type", "text/plain; version=0.0.4")
response.WriteHeader(http.StatusOK)
_, _ = fmt.Fprintf(response, "rvbox_health_ready %d\nrvbox_health_dirty %d\n", boolMetric(ready), boolMetric(dirty))
keys := make([]string, 0, len(counters))
for key := range counters {
_, _ = fmt.Fprintf(response, "rvbox_health_ready %d\nrvbox_health_dirty %d\n", boolMetric(snapshot.Ready), boolMetric(snapshot.Dirty))
keys := make([]string, 0, len(snapshot.Counters))
for key := range snapshot.Counters {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
_, _ = fmt.Fprintf(response, "rvbox_%s_total %d\n", key, counters[key])
_, _ = fmt.Fprintf(response, "rvbox_%s_total %d\n", key, snapshot.Counters[key])
}
keys = keys[:0]
for key := range snapshot.Gauges {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
_, _ = fmt.Fprintf(response, "rvbox_%s %g\n", key, snapshot.Gauges[key])
}
keys = keys[:0]
for key := range snapshot.Histograms {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
histogram := snapshot.Histograms[key]
for index, upperBound := range durationBuckets {
_, _ = fmt.Fprintf(response, "rvbox_%s_seconds_bucket{le=\"%g\"} %d\n", key, upperBound, histogram.Buckets[index])
}
_, _ = fmt.Fprintf(response, "rvbox_%s_seconds_bucket{le=\"+Inf\"} %d\n", key, histogram.Count)
_, _ = fmt.Fprintf(response, "rvbox_%s_seconds_sum %g\nrvbox_%s_seconds_count %d\n", key, histogram.Sum, key, histogram.Count)
}
})
return mux
@@ -142,7 +264,10 @@ func validMetricName(name string) bool {
if name == "" {
return false
}
for _, character := range name {
for index, character := range name {
if index == 0 && !(character == '_' || character >= 'a' && character <= 'z' || character >= 'A' && character <= 'Z') {
return false
}
if !(character == '_' || character >= 'a' && character <= 'z' || character >= 'A' && character <= 'Z' || character >= '0' && character <= '9') {
return false
}
+29
View File
@@ -5,6 +5,7 @@ import (
"net/http/httptest"
"strings"
"testing"
"time"
)
func TestHealthHandlerReadinessAndMetrics_HP_OPS_01(t *testing.T) {
@@ -35,8 +36,36 @@ func TestHealthMetricNamesAreBounded_BH_OPS_01(t *testing.T) {
t.Parallel()
health := New()
health.Inc("bad name")
health.Inc("1bad_name")
_, _, counters := health.Snapshot()
if len(counters) != 0 {
t.Fatalf("invalid metric name recorded: %#v", counters)
}
}
func TestHealthRendersBoundedGaugesAndDurationHistograms_HP_OPS_02(t *testing.T) {
t.Parallel()
health := New()
health.SetGauge("client_spool_bytes", 4096)
health.SetGauge("bad gauge", 1)
health.ObserveDuration("dispatch_latency", 12*time.Millisecond)
health.ObserveDuration("bad latency", time.Second)
metrics := httptest.NewRecorder()
health.Handler(Paths{}).ServeHTTP(metrics, httptest.NewRequest(http.MethodGet, "/metrics", nil))
body := metrics.Body.String()
for _, want := range []string{
"rvbox_client_spool_bytes 4096",
"rvbox_dispatch_latency_seconds_bucket{le=\"0.05\"} 1",
"rvbox_dispatch_latency_seconds_bucket{le=\"+Inf\"} 1",
"rvbox_dispatch_latency_seconds_sum 0.012",
"rvbox_dispatch_latency_seconds_count 1",
} {
if !strings.Contains(body, want) {
t.Fatalf("metrics missing %q:\n%s", want, body)
}
}
if strings.Contains(body, "bad gauge") || strings.Contains(body, "bad latency") {
t.Fatalf("unbounded metric name rendered: %s", body)
}
}
+14 -1
View File
@@ -18,6 +18,7 @@ import (
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
"github.com/rvbox/rvbox/internal/agentproto"
"github.com/rvbox/rvbox/internal/domain"
"github.com/rvbox/rvbox/internal/observability"
"github.com/rvbox/rvbox/internal/server/store"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
@@ -48,6 +49,8 @@ type Options struct {
Now func() time.Time
CursorKey []byte
WakeClient func(string) bool
// Metrics is optional and receives only static metric names.
Metrics *observability.Health
}
type Service struct {
@@ -59,6 +62,7 @@ type Service struct {
now func() time.Time
cursors *domain.CursorCodec
wakeClient func(string) bool
metrics *observability.Health
}
func NewService(options Options) (*Service, error) {
@@ -94,7 +98,7 @@ func NewService(options Options) (*Service, error) {
if err != nil {
return nil, err
}
return &Service{store: options.Store, defaultQueueTTL: options.DefaultQueueTTL, takeoverTTL: options.TakeoverTTL, limits: options.Limits, now: options.Now, cursors: codec, wakeClient: options.WakeClient}, nil
return &Service{store: options.Store, defaultQueueTTL: options.DefaultQueueTTL, takeoverTTL: options.TakeoverTTL, limits: options.Limits, now: options.Now, cursors: codec, wakeClient: options.WakeClient, metrics: options.Metrics}, nil
}
func (service *Service) ListClients(ctx context.Context, request *rvboxv1.ListClientsRequest) (*rvboxv1.ListClientsResponse, error) {
@@ -252,6 +256,7 @@ func (service *Service) RunCommand(ctx context.Context, request *rvboxv1.RunComm
descriptor := spec.GetScript()
digest := sha256.Sum256(request.GetScriptContent())
if descriptor == nil || uint64(len(request.GetScriptContent())) != descriptor.GetSizeBytes() || !bytes.Equal(digest[:], descriptor.GetSha256()) {
service.incMetric("script_verification_failure")
return nil, controlError(codes.InvalidArgument, rvboxv1.ControlError_INVALID_ARGUMENT, "script_content does not match the script descriptor")
}
}
@@ -676,11 +681,19 @@ func (service *Service) AuthorizeClientTakeover(ctx context.Context, request *rv
expires := service.now().UTC().Add(service.takeoverTTL)
value, _, err := service.store.AuthorizeClientTakeover(ctx, store.TakeoverAuthorization{ClientID: request.GetClientId(), ClientInstanceID: [16]byte(instance), RequestID: [16]byte(requestID), ExpiresAt: expires})
if err != nil {
service.incMetric("client_takeover_rejected")
return nil, mapStoreError(err)
}
service.incMetric("client_takeover_authorized")
return &rvboxv1.AuthorizeClientTakeoverResponse{ExpiresAt: timestamppb.New(value)}, nil
}
func (service *Service) incMetric(name string) {
if service != nil && service.metrics != nil {
service.metrics.Inc(name)
}
}
func incidentRecord(view store.IncidentView) *rvboxv1.StorageIncident {
result := &rvboxv1.StorageIncident{IncidentId: domain.UUID(view.IncidentUUID).String(), DetectedAt: timestamppb.New(view.DetectedAt), State: rvboxv1.StorageIncidentState(view.State), Scope: string(view.Scope), ClientId: view.ClientID, Summary: view.Summary, DataLoss: view.DataLoss, AutomaticallyRepairable: view.AutomaticallyRepairable}
if view.ResolvedAt != nil {
+7
View File
@@ -11,6 +11,7 @@ import (
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
"github.com/rvbox/rvbox/internal/agentproto"
"github.com/rvbox/rvbox/internal/domain"
"github.com/rvbox/rvbox/internal/observability"
"github.com/rvbox/rvbox/internal/server/store"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
@@ -70,6 +71,8 @@ func TestControlListAndGetViews_HP_CONTROL_01(t *testing.T) {
func TestRunCommandRequestIDIdempotencyAndValidation_BH_CONTROL_02(t *testing.T) {
service, persistence := newTestService(t)
defer persistence.Close()
metrics := observability.New()
service.metrics = metrics
ctx := context.Background()
registerControlClient(t, persistence, "win-a", rvboxv1.Platform_PLATFORM_WINDOWS, rvboxv1.ShellType_SHELL_POWERSHELL, 3)
issue := fixedIssue(0xa2)
@@ -102,6 +105,10 @@ func TestRunCommandRequestIDIdempotencyAndValidation_BH_CONTROL_02(t *testing.T)
if _, err := service.RunCommand(ctx, badScript); status.Code(err) != codes.InvalidArgument {
t.Fatalf("mismatched script code = %v", status.Code(err))
}
_, _, counters := metrics.Snapshot()
if counters["script_verification_failure"] != 1 {
t.Fatalf("script verification metrics = %#v", counters)
}
badTTL := proto.Clone(request).(*rvboxv1.RunCommandRequest)
badTTL.RequestId = fixedIssue(0xa4).String()
badTTL.QueueTtl = durationpb.New(-time.Second)
+40
View File
@@ -16,6 +16,7 @@ import (
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
"github.com/rvbox/rvbox/internal/agentproto"
"github.com/rvbox/rvbox/internal/domain"
"github.com/rvbox/rvbox/internal/observability"
"github.com/rvbox/rvbox/internal/server/store"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/timestamppb"
@@ -45,6 +46,9 @@ type AgentServer struct {
HeartbeatIdle time.Duration
LivenessTimeout time.Duration
Now func() time.Time
// Metrics is optional so transport tests and embeddings do not need an
// observability endpoint. Every emitted name is a static bounded value.
Metrics *observability.Health
}
func (server *AgentServer) ServeHTTP(response http.ResponseWriter, request *http.Request) {
@@ -53,6 +57,7 @@ func (server *AgentServer) ServeHTTP(response http.ResponseWriter, request *http
return
}
if request.Header.Get("Origin") != "" {
server.incMetric("protocol_error")
http.Error(response, ErrUnexpectedOrigin.Error(), http.StatusForbidden)
return
}
@@ -83,6 +88,7 @@ func (server *AgentServer) ServeHTTP(response http.ResponseWriter, request *http
if err != nil {
return
}
server.incMetric("agent_connection")
defer connection.CloseNow()
server.serveConnection(request.Context(), connection, heartbeat, started)
}
@@ -98,6 +104,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
}
helloEnvelope, err := agentproto.DecodeEnvelope(payload, server.limits(), rvboxv1.Platform_PLATFORM_UNSPECIFIED)
if err != nil || helloEnvelope.GetClientHello() == nil {
server.incMetric("protocol_error")
server.close(connection, websocket.StatusPolicyViolation, "invalid ClientHello")
return
}
@@ -108,11 +115,13 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
reconcileBoundary := server.now()
instanceID, err := domain.ParseUUIDv7(hello.GetClientInstanceId())
if err != nil {
server.incMetric("protocol_error")
server.close(connection, websocket.StatusPolicyViolation, "invalid client instance ID")
return
}
selected, err := domain.SelectProtocol(server.protocolRange(), hello.GetSupportedProtocol())
if err != nil {
server.incMetric("protocol_error")
server.close(connection, websocket.StatusPolicyViolation, "unsupported protocol")
return
}
@@ -126,15 +135,22 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
server.close(connection, websocket.StatusInternalError, "could not persist client capabilities")
return
}
registrationStarted := time.Now()
registration, err := server.Store.RegisterClientSession(parent, store.ClientRegistration{
ClientID: hello.GetClientId(), Platform: uint32(hello.GetPlatform()), Architecture: hello.GetArchitecture(),
DaemonVersion: hello.GetDaemonVersion(), DaemonCWD: hello.GetDaemonCwd(), SupportedShells: shells,
ClientInstanceID: [16]byte(instanceID), SessionID: sessionID, ConnectedAt: server.now(),
})
server.observeMetric("sqlite_write_latency", time.Since(registrationStarted))
if err != nil {
server.incMetric("client_registration_failure")
if errors.Is(err, store.ErrTakeoverRequired) {
server.incMetric("client_takeover_required")
}
server.close(connection, websocket.StatusPolicyViolation, sessionCloseReason(err))
return
}
server.incMetric("client_registration")
registry := server.Registry
handle, err := registry.Install(hello.GetClientId(), sessionID, registration.Generation)
if err != nil {
@@ -227,16 +243,19 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
// to determine liveness.
heartbeat.Observe(time.Since(started))
if messageType != websocket.MessageBinary {
server.incMetric("protocol_error")
server.close(connection, websocket.StatusUnsupportedData, ErrUnexpectedMessage.Error())
return
}
envelope, err := agentproto.DecodeEnvelope(payload, server.limits(), hello.GetPlatform())
if err != nil || envelope.GetClientHello() != nil || envelope.GetSessionId() != encodedSessionID || envelope.GetSessionGeneration() != registration.Generation {
server.incMetric("stale_message")
server.close(connection, websocket.StatusPolicyViolation, ErrStaleSession.Error())
return
}
clientID, err := server.Store.ValidateLiveSession(sessionContext, sessionID, registration.Generation)
if err != nil || clientID != hello.GetClientId() {
server.incMetric("stale_message")
server.close(connection, websocket.StatusPolicyViolation, ErrStaleSession.Error())
return
}
@@ -317,14 +336,20 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
}
appendEvent, eventErr := eventAppendFromWire(event, hello.GetClientId(), registration.Generation, server.now())
if eventErr != nil {
server.incMetric("protocol_error")
server.close(connection, websocket.StatusPolicyViolation, "invalid command event")
return
}
appended, eventErr := server.Store.AppendCommandEvent(sessionContext, appendEvent)
if eventErr != nil {
server.incMetric("command_event_rejected")
server.close(connection, websocket.StatusPolicyViolation, "command event was not accepted")
return
}
server.incMetric("command_event")
if event.GetLifecycle() != nil {
server.incMetric("command_transition")
}
ack, eventErr := proto.Marshal(&rvboxv1.AgentEnvelope{SessionId: encodedSessionID, SessionGeneration: registration.Generation, Payload: &rvboxv1.AgentEnvelope_EventAck{EventAck: &rvboxv1.EventAck{IssueUuid: event.GetIssueUuid(), ThroughEventSeq: appended.ThroughEventSeq}}})
if eventErr != nil || queue.EnqueueControl(Frame{Kind: FrameControl, Payload: ack}) != nil {
server.close(connection, websocket.StatusInternalError, "could not acknowledge command event")
@@ -452,6 +477,8 @@ func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *Write
beforeEnqueue(candidate)
}
if !queue.EnqueueData(Frame{Kind: FrameData, Payload: encoded, OnWritten: func() {
server.incMetric("command_dispatch")
server.observeMetric("dispatch_latency", server.now().Sub(candidate.IssueTime))
if onWritten != nil {
onWritten(candidate.IssueUUID)
}
@@ -756,11 +783,24 @@ func (server *AgentServer) writeLoop(ctx context.Context, connection *websocket.
return
}
case HeartbeatClose:
server.incMetric("heartbeat_timeout")
return
}
}
}
func (server *AgentServer) incMetric(name string) {
if server != nil && server.Metrics != nil {
server.Metrics.Inc(name)
}
}
func (server *AgentServer) observeMetric(name string, duration time.Duration) {
if server != nil && server.Metrics != nil {
server.Metrics.ObserveDuration(name, duration)
}
}
func (server *AgentServer) close(connection *websocket.Conn, status websocket.StatusCode, reason string) {
_ = connection.Close(status, reason)
}
@@ -18,6 +18,7 @@ import (
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
"github.com/rvbox/rvbox/internal/agentproto"
"github.com/rvbox/rvbox/internal/domain"
"github.com/rvbox/rvbox/internal/observability"
"github.com/rvbox/rvbox/internal/server/store"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/timestamppb"
@@ -98,6 +99,19 @@ func TestAgentServerRegistrationAndReplacement_HP_SES_05(t *testing.T) {
}
}
func TestAgentServerEmitsBoundedRegistrationMetrics_HP_OPS_05(t *testing.T) {
server, agent, cleanup := newTestAgentServerWithStore(t)
defer cleanup()
metrics := observability.New()
agent.Metrics = metrics
connection, _ := dialAndHello(t, server, "metric-client", "019c46f1-1d02-7000-8000-000000000073")
defer connection.CloseNow()
_, _, counters := metrics.Snapshot()
if counters["agent_connection"] != 1 || counters["client_registration"] != 1 {
t.Fatalf("registration metrics = %#v", counters)
}
}
func TestAgentServerDispatchesCommandAdmittedDuringReconcile_HP_SES_11(t *testing.T) {
server, agent, cleanup := newTestAgentServerWithStore(t)
defer cleanup()
+34
View File
@@ -0,0 +1,34 @@
package store
import (
"context"
"errors"
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
)
// Telemetry is the bounded operational state that the server publishes as
// gauges. It deliberately contains no client or command identifier.
type Telemetry struct {
QueuedCommands uint64
ChargedBytes uint64
}
func (store *Store) Telemetry(ctx context.Context) (Telemetry, error) {
if store == nil {
return Telemetry{}, errors.New("server store is unavailable")
}
store.mu.Lock()
defer store.mu.Unlock()
if store.db == nil {
return Telemetry{}, ErrStoreClosed
}
var result Telemetry
if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM commands WHERE lifecycle = ?`, uint32(rvboxv1.CommandLifecycle_COMMAND_QUEUED)).Scan(&result.QueuedCommands); err != nil {
return Telemetry{}, err
}
if err := store.db.QueryRowContext(ctx, `SELECT command_charged_bytes FROM storage_counters WHERE singleton = 1`).Scan(&result.ChargedBytes); err != nil {
return Telemetry{}, err
}
return result, nil
}
+35
View File
@@ -0,0 +1,35 @@
package store
import (
"context"
"crypto/sha256"
"path/filepath"
"testing"
"time"
"github.com/rvbox/rvbox/internal/domain"
)
func TestTelemetryReportsQueuedDurableState_HP_OPS_06(t *testing.T) {
t.Parallel()
ctx := context.Background()
persistence, err := Open(ctx, Options{DataDir: filepath.Join(t.TempDir(), "state"), BusyTimeout: time.Second})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = persistence.Close() })
if _, err := persistence.RegisterClientSession(ctx, ClientRegistration{ClientID: "telemetry-client", Platform: 2, Architecture: "amd64", DaemonVersion: "test", DaemonCWD: `C:\`, SupportedShells: []byte{1}, ClientInstanceID: [16]byte{1}, SessionID: [16]byte{2}, ConnectedAt: time.Now().UTC()}); err != nil {
t.Fatal(err)
}
issue, err := domain.ParseUUIDv7("019c46f1-1d02-7000-8000-0000000000d1")
if err != nil {
t.Fatal(err)
}
if _, err := persistence.QueueCommand(ctx, QueueCommandInput{IssueUUID: issue, ClientID: "telemetry-client", IssueTime: time.Now().UTC(), ReceiptTime: time.Now().UTC(), ImmutableSHA256: sha256.Sum256([]byte("telemetry")), ExecutionSpec: []byte("spec")}); err != nil {
t.Fatal(err)
}
telemetry, err := persistence.Telemetry(ctx)
if err != nil || telemetry.QueuedCommands != 1 || telemetry.ChargedBytes == 0 {
t.Fatalf("telemetry = %#v, %v", telemetry, err)
}
}
+143
View File
@@ -0,0 +1,143 @@
#!/bin/sh
# Build an inspectable, reproducible RVBox Linux-server/Windows-client bundle.
# Publishing and certificate custody remain outside this repository; an optional
# local signing hook can sign the Windows executable before checksums are made.
set -eu
repo_root=$(CDPATH= cd -- "$(dirname -- "$0")/.." && pwd)
compose_file=$repo_root/deploy/compose.yaml
usage() {
cat <<'EOF'
usage: scripts/release build --version VERSION --bootstrap-server-url WSS_URL [--output DIR] [--sign-windows-hook FILE]
Builds a fresh immutable bundle containing:
rvbox-server-linux-amd64
rvbox-server-linux-arm64
rvc-linux-amd64
rvc-linux-arm64
rvbox-windows-amd64.exe
SHA256SUMS
manifest.json
The default output is dist/rvbox-VERSION. --output must remain below this
repository's ignored dist/ tree. VERSION must be a safe release label
([0-9A-Za-z][0-9A-Za-z._+-]{0,63}); an existing final output is never replaced.
WSS_URL is the deliberate first-install endpoint embedded only in the Windows
installer; it must be an absolute `wss://` URL and is never inferred from a
test fixture. Artifacts are built inside the pinned Docker toolchain with that
exact version embedded in `--version`. The optional signing hook is a regular executable run
on the host as: HOOK WINDOWS_EXE VERSION. It receives RVBOX_ARTIFACT and
RVBOX_VERSION as environment variables and must sign the supplied executable
in place. SHA256SUMS and manifest.json are generated only after it succeeds.
No signing hook means the manifest explicitly marks the Windows artifact as
unsigned. This is deliberate for CI/test builds; do not publish it as signed.
EOF
}
fail() { printf '%s\n' "release: $*" >&2; exit 2; }
command=${1:-}
[ "$#" -gt 0 ] && shift
[ "$command" = build ] || { usage >&2; fail "expected build"; }
version=
output=
sign_hook=
bootstrap_server_url=
while [ "$#" -gt 0 ]; do
case $1 in
--version) [ "$#" -ge 2 ] || fail "--version needs a value"; version=$2; shift 2 ;;
--bootstrap-server-url) [ "$#" -ge 2 ] || fail "--bootstrap-server-url needs a value"; bootstrap_server_url=$2; shift 2 ;;
--output) [ "$#" -ge 2 ] || fail "--output needs a directory"; output=$2; shift 2 ;;
--sign-windows-hook) [ "$#" -ge 2 ] || fail "--sign-windows-hook needs an executable file"; sign_hook=$2; shift 2 ;;
--help|-h) usage; exit 0 ;;
*) fail "unknown argument $1" ;;
esac
done
case $version in
''|[!0-9A-Za-z]*|*[!0-9A-Za-z._+-]*|?????????????????????????????????????????????????????????????????*)
fail "--version must match [0-9A-Za-z][0-9A-Za-z._+-]{0,63}"
;;
esac
case $bootstrap_server_url in
wss://*) ;;
*) fail "--bootstrap-server-url must be an absolute wss:// URL" ;;
esac
case $bootstrap_server_url in
*[!A-Za-z0-9:/._-]*) fail "--bootstrap-server-url contains unsupported characters" ;;
esac
if [ -z "$output" ]; then output=$repo_root/dist/rvbox-$version; fi
case $output in
/*) ;;
*) output=$repo_root/$output ;;
esac
case $output in
"$repo_root"/dist/*) ;;
*) fail "--output must be below $repo_root/dist" ;;
esac
parent=$(dirname -- "$output")
[ ! -e "$output" ] || fail "refusing to replace existing output $output"
[ -d "$parent" ] || mkdir -p "$parent"
[ ! -L "$parent" ] || fail "refusing symlink output parent $parent"
if [ -n "$sign_hook" ]; then
[ -f "$sign_hook" ] && [ ! -L "$sign_hook" ] && [ -x "$sign_hook" ] || fail "signing hook must be an executable regular file"
fi
docker version >/dev/null 2>&1 || fail "Docker is unavailable"
docker compose -f "$compose_file" version >/dev/null 2>&1 || fail "Docker Compose is unavailable"
partial=$parent/.rvbox-$version.partial-$$
mkdir "$partial" || fail "could not create private release staging directory"
cleanup() { rm -rf -- "$partial"; }
trap cleanup EXIT HUP INT TERM
container_partial=/workspace/${partial#"$repo_root"/}
commit=$(git -C "$repo_root" rev-parse HEAD)
dirty=$(git -C "$repo_root" status --porcelain)
[ -z "$dirty" ] || fail "refusing release build from a dirty worktree"
docker compose -f "$compose_file" run --rm toolchain sh -ec '
set -eu
version=$1
out=$2
flags="-s -w -X main.buildVersion=$version"
windows_flags="$flags -X main.bootstrapServerURL=$3"
resource=/workspace/cmd/rvbox/rvbox-release_windows_amd64.syso
icon="$out/.rvbox-release-icon.png"
trap "rm -f -- \"$resource\" \"$icon\"" EXIT HUP INT TERM
go run ./tools/releaseicon --out "$icon"
go-winres simply --arch amd64 --out /workspace/cmd/rvbox/rvbox-release --icon "$icon" --manifest gui --file-version "$version" --product-version "$version" --file-description "RVBox Windows client" --product-name "RVBox" --original-filename "rvbox.exe"
CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -trimpath -ldflags "$flags" -o "$out/rvbox-server-linux-amd64" ./cmd/rvbox-server
CGO_ENABLED=0 GOOS=linux GOARCH=arm64 go build -trimpath -ldflags "$flags" -o "$out/rvbox-server-linux-arm64" ./cmd/rvbox-server
CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -trimpath -ldflags "$flags" -o "$out/rvc-linux-amd64" ./cmd/rvc
CGO_ENABLED=0 GOOS=linux GOARCH=arm64 go build -trimpath -ldflags "$flags" -o "$out/rvc-linux-arm64" ./cmd/rvc
CGO_ENABLED=0 GOOS=windows GOARCH=amd64 go build -trimpath -ldflags "-H=windowsgui $windows_flags" -o "$out/rvbox-windows-amd64.exe" ./cmd/rvbox
' sh "$version" "$container_partial" "$bootstrap_server_url"
signed=false
if [ -n "$sign_hook" ]; then
RVBOX_ARTIFACT=$partial/rvbox-windows-amd64.exe RVBOX_VERSION=$version "$sign_hook" "$partial/rvbox-windows-amd64.exe" "$version"
signed=true
fi
for artifact in rvbox-server-linux-amd64 rvbox-server-linux-arm64 rvc-linux-amd64 rvc-linux-arm64 rvbox-windows-amd64.exe; do
[ -f "$partial/$artifact" ] && [ ! -L "$partial/$artifact" ] || fail "builder did not create regular $artifact"
done
(cd "$partial" && sha256sum rvbox-server-linux-amd64 rvbox-server-linux-arm64 rvc-linux-amd64 rvc-linux-arm64 rvbox-windows-amd64.exe) >"$partial/SHA256SUMS"
{
printf '{\n'
printf ' "schema": 1,\n'
printf ' "version": "%s",\n' "$version"
printf ' "git_commit": "%s",\n' "$commit"
printf ' "windows_signed": %s,\n' "$signed"
printf ' "artifacts": "SHA256SUMS"\n'
printf '}\n'
} >"$partial/manifest.json"
chmod 0755 "$partial/rvbox-server-linux-amd64" "$partial/rvbox-server-linux-arm64" "$partial/rvc-linux-amd64" "$partial/rvc-linux-arm64" "$partial/rvbox-windows-amd64.exe"
chmod 0644 "$partial/SHA256SUMS" "$partial/manifest.json"
mv -- "$partial" "$output"
trap - EXIT HUP INT TERM
printf 'release_bundle=%s\n' "$output"
+16 -2
View File
@@ -20,7 +20,7 @@ Actions:
stage copy a bundle containing rvbox.exe and client.toml into the guest test root
install install and start RVBox from the staged bundle through a fixture-only full-admin principal
run start the already-installed RVBox SCM service from the staged bundle
logoff log off the sole active fixture user; use only after service installation
logoff log off the sole active fixture user; use only after service installation
collect copy bounded guest artifacts to the local test-run directory
inspect read-only RVBox SCM state and bounded client log from a prepared run
logs read-only service-startup and client logs from a prepared run
@@ -82,7 +82,7 @@ while [ "$#" -gt 0 ]; do
--run-id) [ "$#" -ge 2 ] || fail "--run-id needs a value"; run_id=$2; shift 2 ;;
--bundle) [ "$#" -ge 2 ] || fail "--bundle needs a value"; bundle=$2; shift 2 ;;
--endpoint) [ "$#" -ge 2 ] || fail "--endpoint needs a value"; endpoint=$2; shift 2 ;;
--fail-contexts) [ "$#" -ge 2 ] || fail "--fail-contexts needs a value"; fail_contexts=$2; shift 2 ;;
--fail-contexts) [ "$#" -ge 2 ] || fail "--fail-contexts needs a value"; fail_contexts=$2; shift 2 ;;
--help|-h) usage; exit 0 ;;
*) fail "unknown argument $1" ;;
esac
@@ -316,6 +316,18 @@ wait_guest_additions() {
fail "Guest Additions did not publish version and Windows OS-release properties"
}
configure_display_keepalive() {
# VirtualBox VRDE reports a zero-bpp framebuffer when Windows powers off
# the virtual monitor. That leaves an authenticated Guacamole tunnel in
# its "Waiting for response" state even though RDP negotiation succeeded.
# Keep the disposable diagnostic fixture awake; this is deliberately
# applied after every snapshot restore because power-plan state belongs to
# the guest snapshot, not to the controller.
provisioner_run --exe 'C:\Windows\System32\cmd.exe' --wait-stdout --wait-stderr --unquoted-args -- \
/d /s /c '(powercfg /change monitor-timeout-ac 0 && powercfg /change monitor-timeout-dc 0 && powercfg /change standby-timeout-ac 0 && powercfg /change standby-timeout-dc 0 && powercfg /change hibernate-timeout-ac 0 && powercfg /change hibernate-timeout-dc 0 && echo RVBOX_GUEST_OK)' >/dev/null || \
fail "fixture provisioner could not disable disposable display/sleep timers"
}
wait_service() {
wanted=$1
attempt=0
@@ -366,6 +378,8 @@ case "$action" in
step prepare-vm-started
wait_guest_additions
step prepare-guest-additions-ready
configure_display_keepalive
step prepare-display-keepalive
assert_clean_guest
step prepare-clean-baseline-verified
;;
+52
View File
@@ -17,6 +17,21 @@ layer = "unit"
status = "implemented"
tests = ["internal/config/effective_test.go:TestEffectiveConfigurationGolden_HP_CFG_09"]
[[requirements]]
id = "HP-CFG-10"
layer = "unit"
status = "implemented"
tests = ["internal/config/config_test.go:TestClientMayDisableObservabilityListener_HP_CFG_10"]
[[requirements]]
id = "HP-WINCLI-04"
layer = "unit"
status = "implemented"
tests = [
"cmd/rvbox/main_test.go:TestPackagedClientConfigUsesMinimalSafeDefaults_HP_WINCLI_04",
"cmd/rvbox/main_test.go:TestInstallerClientIDNormalizesHostnames_HP_WINCLI_05",
]
[[requirements]]
id = "HP-CTL-06"
layer = "unit"
@@ -182,6 +197,7 @@ tests = [
"internal/agentproto/validate_test.go:FuzzDecodeOutputChunkBounded_SEC_PROTO_02",
"internal/server/control/jsonrpc_test.go:FuzzDecodeJSONRPCRequestBounded_SEC_CTL_01",
"internal/server/store/segment_test.go:FuzzDecodeSegmentRecordDoesNotEscapeBounds",
"internal/domain/cursor_test.go:FuzzDecodeCursorBounded_SEC_CTL_02",
]
[[requirements]]
@@ -637,6 +653,42 @@ layer = "unit"
status = "implemented"
tests = ["internal/observability/rotate_test.go:TestFormatLogWritesStructuredBoundedRecords_HP_OPS_04"]
[[requirements]]
id = "HP-OPS-05"
layer = "unit"
status = "implemented"
tests = [
"internal/observability/health_test.go:TestHealthRendersBoundedGaugesAndDurationHistograms_HP_OPS_02",
"internal/server/session/agent_server_test.go:TestAgentServerEmitsBoundedRegistrationMetrics_HP_OPS_05",
]
[[requirements]]
id = "HP-OPS-06"
layer = "unit"
status = "implemented"
tests = [
"internal/server/store/telemetry_test.go:TestTelemetryReportsQueuedDurableState_HP_OPS_06",
"internal/client/spool/telemetry_test.go:TestTelemetryReportsQueuedDurableSpool_HP_OPS_07",
]
[[requirements]]
id = "HP-OPS-07"
layer = "unit"
status = "implemented"
tests = [
"cmd/rvbox-server/main_test.go:TestPublishServerTelemetrySample_HP_OPS_07",
"cmd/rvbox/main_test.go:TestPublishClientTelemetrySample_HP_OPS_07",
]
[[requirements]]
id = "BH-OPS-04"
layer = "unit"
status = "implemented"
tests = [
"cmd/rvbox-server/main_test.go:TestPublishServerTelemetrySampleFailure_BH_OPS_04",
"cmd/rvbox/main_test.go:TestPublishClientTelemetrySampleFailure_BH_OPS_04",
]
[[requirements]]
id = "HP-DISPATCH-10"
layer = "unit"
+200 -19
View File
@@ -9,10 +9,11 @@ The helper builds this private path:
```text
browser -- HTTPS/self-signed --> nginx + Guacamole containers
(IPv4 and, in public mode, a tracked IPv6 forward)
|
private Docker gateway
|
controller SSH tunnel --> Helium 127.0.0.1:3389 --> VirtualBox VRDE --> VM console
controller SSH watchdog/tunnel --> Helium 127.0.0.1:3389 --> VirtualBox VRDE --> VM console
```
Only the HTTPS listener can be made public, and that requires an explicit
@@ -23,6 +24,63 @@ password to the VM. Clipboard, drives, printing, audio, microphone input, GFX,
and display-resize extensions are disabled because the fixture's VirtualBox
RDP4 server does not handle them reliably.
## Current fixture quick start
The following is the copy/paste path for the documented `rvbox-win10-test`
fixture. Run it from the repository root. The first command acquires the
exclusive VM lease, restores the clean administrator-enabled snapshot, starts
the VM, and applies the VRDE display keepalive. The second command starts the
Docker-only gateway at the currently published endpoint:
```sh
scripts/windows/test-host prepare --run-id interactive-rdp
test/rdp-access/rdp-access up \
--bind 0.0.0.0 \
--public-host x1.xcel.me
```
`up` prompts for the fixture password. If the password is not already in the
operator's secure notes, the test-only value is available only in the Helium
host's mode-600 file documented in [`docs/testing-vm.md`](../../docs/testing-vm.md).
For an unattended first start, pass that file through the existing SSH alias
without putting the value in an argument, environment variable, or log:
```sh
ssh helium-remote 'cat /home/cabbage/.local/share/rvbox-secrets/rvbox-win10-test.password' |
test/rdp-access/rdp-access up \
--bind 0.0.0.0 \
--public-host x1.xcel.me \
--web-password-stdin
```
The browser URL is `https://x1.xcel.me:5002/guacamole/`; accept the temporary
self-signed certificate warning and sign in as `rvboxtest`. The same command
is safe to rerun after a tunnel interruption. If the stack is already running,
`up` reuses it and repairs only missing forwarding; it does not consume stdin.
Use `--reset-auth` only after `down` when intentionally replacing the saved
verifier, for example:
```sh
test/rdp-access/rdp-access down
ssh helium-remote 'cat /home/cabbage/.local/share/rvbox-secrets/rvbox-win10-test.password' |
test/rdp-access/rdp-access up \
--bind 0.0.0.0 \
--public-host x1.xcel.me \
--reset-auth \
--web-password-stdin
```
Do not run `clean` or `test-host reset` while another agent owns the fixture
lease or has a live browser session. The generated verifier, certificate,
tunnel state, and logs are disposable and live only under the ignored
`test/rdp-access/.runtime/` directory.
Before a first `prepare`, run `scripts/windows/test-host status`. If it reports
`state=running`, another run may already own the fixture lease; do not restore
or reset it. Reuse that run's access path or coordinate with its owner. If it
reports `state=poweroff`, the quick-start sequence above is safe. The helper
does not guess at ownership or silently take over a running VM.
## Lifecycle
Run commands from the repository root. First prepare the disposable VM using a
@@ -53,18 +111,88 @@ For localhost-only use, omit `--bind` and `--public-host`. The default listener
is `127.0.0.1:5002`; use a local SSH forward or a browser on the controller.
Choose alternate ports with `--http-port` and `--tunnel-port` if either is in
use. The VM must already be running. `up` checks the documented VM/snapshot
identity but intentionally does not restore, start, stop, or reset the VM.
identity but intentionally does not restore, start, stop, or reset the VM. It
starts a small host-side watchdog for the SSH master. Before printing the URL,
`up` waits until the HTTPS gateway can serve `/guacamole/`; this prevents a
browser login from racing Tomcat's Guacamole WAR deployment. The watchdog
reconnects after a transient Helium/SSH failure while preserving the same
Docker-gateway listener, so an already-open Guacamole session can recover
without restarting the containers. Its PID, stop marker, and diagnostic log
are kept under the ignored `.runtime/` directory.
All fixture-specific values have embedded, working defaults: the
`helium-remote` SSH alias, Helium's loopback VRDE endpoint (`127.0.0.1:3389`),
`rvboxtest`, the browser listener (`127.0.0.1:5002`), and the private tunnel
port (`54001`). They can be overridden without editing tracked files through
`RVBOX_TEST_VBOX_HOST`, `RDP_ACCESS_VRDE_HOST`, `RDP_ACCESS_VRDE_PORT`,
`RDP_ACCESS_PROFILE`,
`RDP_ACCESS_WEB_USER`, `RDP_ACCESS_RDP_USER`, `RDP_ACCESS_BIND`,
`RDP_ACCESS_HTTP_PORT`, `RDP_ACCESS_TUNNEL_PORT`, and
`RDP_ACCESS_PUBLIC_HOST`. The helper deliberately limits the VRDE host to
`RDP_ACCESS_PUBLIC_HOST`, `RDP_ACCESS_IPV6_BIND`, and
`RDP_ACCESS_IPV6_FORWARD`. The helper deliberately limits the VRDE host to
Helium loopback (`127.0.0.1` or `localhost`) so an override cannot accidentally
turn the diagnostic server into a remote target.
Set `RDP_ACCESS_GUACD_LOG_LEVEL=debug` temporarily when collecting detailed
guacd/RDP negotiation diagnostics; the default is `info`.
The VM selector and the gateway profile are separate. Set the matching
`RVBOX_TEST_VBOX_*` identity variables (and `RDP_ACCESS_VRDE_PORT`) for the
VM you want; the complete identity set is listed in
[`docs/testing-vm.md`](../../docs/testing-vm.md). For two already-running VMs,
give the second gateway a distinct profile and listener ports, for example:
```sh
RDP_ACCESS_PROFILE=vm-b \
RVBOX_TEST_VBOX_HOST=helium-remote-b \
RDP_ACCESS_VRDE_PORT=3391 \
RDP_ACCESS_HTTP_PORT=5003 \
RDP_ACCESS_TUNNEL_PORT=54002 \
test/rdp-access/rdp-access up --bind 0.0.0.0 --public-host vm-b.example.net
```
Replace `helium-remote-b`, `3391`, and the public hostname with the second
fixture's recorded values, and export its matching VM/snapshot UUID variables
before running `up` or `status`. A non-`default` profile stores its verifier,
certificate, tunnel state, and logs under `.runtime/<profile>/` and uses the
Compose project `rvbox-rdp-access-<profile>`. Keep the HTTP port, public
hostname, and profile unique so IPv4/IPv6 listeners and browser sessions do
not collide. Inspect, repair, stop, or clean that instance by repeating the
same `RDP_ACCESS_PROFILE` and endpoint variables on `status`, `repair`,
`down`, or `clean`.
Profile isolation does not bypass the native controller's exclusive lease:
`test-host prepare`/`reset` must still be coordinated when VMs share one
fixture host and staging root. For VMs that are already running, the profile
and matching identity/VRDE variables are sufficient to select the endpoint.
When `--bind 0.0.0.0` is used, `up` also starts a tracked `socat` listener on
`[::]:$RDP_ACCESS_HTTP_PORT` and forwards it to the IPv4 gateway listener. This
matters when the public hostname has an AAAA record: some browsers choose IPv6
for the WebSocket even if the initial page used IPv4. The listener is recorded
under `.runtime/ipv6-forward.pid` and is removed by `down` and `clean`. Its
default bind is `::` (`RDP_ACCESS_IPV6_BIND`); set
`RDP_ACCESS_IPV6_FORWARD=never` to deliberately use IPv4 only, `always` to make
IPv6 startup a hard requirement, or leave the default `auto` for a warning and
an IPv4-only fallback when `socat` is unavailable. `status` reports
`ipv6_forward=active`, `active_external`, `inactive`, or `disabled`.
If `up` is run again while the Compose services are still running, it is
idempotent: an existing healthy tunnel is reused, and a missing tunnel is
recreated in place. `status` reports a helper-owned tunnel as
`private_tunnel=active` and a listener supplied by an external/interactive
SSH supervisor as `private_tunnel=active_external`. A missing listener is
reported as `private_tunnel=inactive`; inspect `.runtime/tunnel.log` and run
`up` again to trigger a bounded reconnect attempt.
If `up` reports that an existing gateway has a different configuration, or a
previous start left only part of the Compose stack, run `down` once and then
repeat `up`. Changing the bind address, browser host, or either port likewise
requires `down` first; refusing to mutate a live stack prevents an old browser
session from silently reaching a different endpoint.
The XML mapping intentionally uses Guacamole's `${GUAC_PASSWORD}` connection
parameter token. Guacamole resolves this to the password entered at web login;
it is not a host environment variable and must remain in the template.
After the interactive action, close the browser connection and remove the
temporary access path before releasing the fixture lease:
@@ -74,9 +202,10 @@ test/rdp-access/rdp-access down
scripts/windows/test-host reset --run-id interactive-rdp
```
`down` stops containers and the SSH master/tunnel but retains the one-day
certificate and password verifier for a quick restart. To remove all generated
state, including the certificate and verifier:
`down` stops containers and the SSH watchdog/master/tunnel but retains the
one-day certificate and password verifier for a quick restart. To remove all
generated state, including the certificate, verifier, watchdog PID, and tunnel
log:
```sh
test/rdp-access/rdp-access clean
@@ -101,18 +230,70 @@ test/rdp-access/rdp-access logs --tail=100
test/rdp-access/rdp-access url
```
If the browser reaches Guacamole but stays on “Waiting for response”, verify
that the VM is running and the private tunnel is active with `status`. This
helper already uses `security=rdp` and disables Guacamole's GFX extension,
which are required by the fixture's legacy VRDE server. Do not switch the
helper to native Windows RDP: `TermService` is intentionally disabled in the
baseline. If VRDE remains unusable, stop this helper and use Guest Control for
the deterministic portion of the work; record the blocked interactive step in
the native test report.
If the browser reaches Guacamole but stays on “Waiting for response”, run
`status` first. Confirm `private_tunnel=active` (or
`active_external`) and, for a public dual-stack hostname,
`ipv6_forward=active` (or a known-good external forward). Then check the VM
state and the last lines of
`.runtime/tunnel.log`. A tunnel can be recreated without losing the Compose
stack by running `up` again. This helper already uses `security=rdp` and
disables Guacamole's GFX extension, which are required by the fixture's legacy
VRDE server. Do not switch the helper to native Windows RDP:
`TermService` is intentionally disabled in the baseline. If VRDE remains
unusable, stop this helper and use Guest Control for the deterministic portion
of the work; record the blocked interactive step in the native test report.
The helper requires Docker/Docker Compose, SSH access through the existing
`helium-remote` alias, and the fixture password file documented in
If guacd has a stale worker and the VM reports that multiple connections are
disabled, run:
```sh
test/rdp-access/rdp-access repair
```
`repair` restarts only guacd, retains the VM, SSH tunnel, gateway, and IPv6
forward, and does not alter the Windows guest. Close old browser tabs and sign
in again after it completes.
For quick recovery, use this decision table:
| Need | Command | Effect |
| --- | --- | --- |
| Inspect everything | `test/rdp-access/rdp-access status` | Read-only VM, Compose, tunnel, and IPv6-forward state |
| Start or reconnect access | `test/rdp-access/rdp-access up --bind 0.0.0.0 --public-host x1.xcel.me` | Reuses healthy services; recreates only missing tunnel/forward |
| Browser says “Waiting for response” | `test/rdp-access/rdp-access status`; `test/rdp-access/rdp-access logs --tail=100`; if guacd is stale, `test/rdp-access/rdp-access repair` | Diagnoses tunnel/dual-stack issues; restarts only guacd |
| Partial stack or changed endpoint | `test/rdp-access/rdp-access down`, then `test/rdp-access/rdp-access up --bind 0.0.0.0 --public-host x1.xcel.me` | Recreates the exact project with the new settings |
| Print the current URL | `test/rdp-access/rdp-access url` | Read-only URL from the saved session |
| Stop temporary access | `test/rdp-access/rdp-access down` | Stops gateway, IPv6 forward, SSH watchdog/tunnel; retains verifier/cert |
| Reclaim generated state | `test/rdp-access/rdp-access clean` | Stops access and removes only `.runtime/` |
| Reclaim exact unused images too | `test/rdp-access/rdp-access clean --images` | Also removes owned Guacamole/nginx images; leaves shared Alpine intact |
| Return VM to clean baseline | `scripts/windows/test-host reset --run-id interactive-rdp` | Stops/ restores the exact leased VM snapshot and leaves it powered off |
The two scripts have a deliberate ownership boundary. `rdp-access` manages only
the temporary browser gateway, Guacamole/guacd containers, IPv6 forward, and
Helium SSH tunnel. `scripts/windows/test-host` manages the leased VM and its
service lifecycle. For a prepared run, the remaining VM-side management actions
are:
```sh
scripts/windows/test-host inspect --run-id interactive-rdp # SCM/read-only client state
scripts/windows/test-host logs --run-id interactive-rdp # bounded startup/client logs
scripts/windows/test-host collect --run-id interactive-rdp # bounded artifacts
scripts/windows/test-host stop --run-id interactive-rdp # stop service and request guest shutdown
scripts/windows/test-host recover --run-id interactive-rdp # inspect an interrupted stopped run
```
`stage`, `install`, `run`, and `logoff` are also available in the
`test-host --help` command reference for a complete native-service run; they
are not required merely to use the RDP gateway. `native-test recover` and
`native-test clean` are the corresponding whole-lane recovery/cleanup actions
when the Linux server stack is part of the run.
The helper requires Docker/Docker Compose with Docker socket access for the
invoking account (`docker version` must succeed), `socat` for the optional
public dual-stack forward, SSH access through the existing `helium-remote`
alias, and the fixture password file documented in
[`docs/testing-vm.md`](../../docs/testing-vm.md). It does not install host
packages, write credentials into Git, or alter VM settings. The tracked files
are Docker-only configuration and the controller script; all generated content
is ignored beneath `.runtime/`.
packages, write credentials into Git, or alter VM identity/settings. The
authoritative `test-host prepare` step applies only the disposable
display/sleep keepalive needed by VirtualBox VRDE; all generated content is
ignored beneath `.runtime/`.
+2
View File
@@ -1,6 +1,8 @@
services:
guacd:
image: guacamole/guacd:1.6.0
environment:
LOG_LEVEL: ${RDP_ACCESS_GUACD_LOG_LEVEL:-info}
guacamole:
image: guacamole/guacamole:1.6.0
+3
View File
@@ -15,5 +15,8 @@ server {
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_buffering off;
proxy_read_timeout 1h;
proxy_send_timeout 1h;
proxy_set_header X-Accel-Buffering no;
}
}
+347 -21
View File
@@ -5,17 +5,20 @@ set -eu
helper_dir=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd)
repo_root=$(CDPATH= cd -- "$helper_dir/../.." && pwd)
runtime_dir=$helper_dir/.runtime
runtime_root=$helper_dir/.runtime
compose_file=$helper_dir/compose.yaml
mapping_template=$helper_dir/user-mapping.xml.in
project=rvbox-rdp-access
: "${RDP_ACCESS_PROFILE:=default}"
: "${RDP_ACCESS_BIND:=127.0.0.1}"
: "${RDP_ACCESS_HTTP_PORT:=5002}"
: "${RDP_ACCESS_TUNNEL_PORT:=54001}"
: "${RDP_ACCESS_PUBLIC_HOST:=localhost}"
: "${RDP_ACCESS_WEB_USER:=rvboxtest}"
: "${RDP_ACCESS_RDP_USER:=rvboxtest}"
: "${RDP_ACCESS_GUACD_LOG_LEVEL:=info}"
: "${RDP_ACCESS_IPV6_BIND:=::}"
: "${RDP_ACCESS_IPV6_FORWARD:=auto}"
: "${RVBOX_TEST_VBOX_HOST:=helium-remote}"
: "${RDP_ACCESS_VRDE_HOST:=127.0.0.1}"
: "${RDP_ACCESS_VRDE_PORT:=3389}"
@@ -27,13 +30,16 @@ usage: test/rdp-access/rdp-access ACTION [OPTIONS]
Actions:
up start a private VRDE SSH tunnel and self-signed HTTPS Guacamole
status show gateway, tunnel, and fixture status without changing anything
repair restart only guacd to release a stale single-VRDE-slot worker
url print the current browser URL
logs follow or print Compose logs (pass Docker Compose log options)
down stop containers and the private SSH tunnel; retain generated state
down stop containers and the private SSH watchdog/tunnel; retain generated state
clean run down and delete generated state; pass --images to also remove
the exact unused Guacamole/nginx images
up options:
--profile NAME isolated VM/gateway profile (default default; use one
distinct profile per simultaneous VM gateway)
--bind ADDRESS listener address (default 127.0.0.1; use 0.0.0.0 only
for a temporary, deliberately public endpoint)
--http-port PORT HTTPS listener port (default 5002)
@@ -43,10 +49,13 @@ up options:
--reset-auth discard the saved password hash and prompt again
--web-password-stdin read the password once from stdin instead of prompting
Environment equivalents: RDP_ACCESS_BIND, RDP_ACCESS_HTTP_PORT,
Environment equivalents: RDP_ACCESS_PROFILE, RDP_ACCESS_BIND, RDP_ACCESS_HTTP_PORT,
RDP_ACCESS_TUNNEL_PORT, RDP_ACCESS_PUBLIC_HOST, RDP_ACCESS_WEB_USER,
RDP_ACCESS_RDP_USER, RVBOX_TEST_VBOX_HOST, RDP_ACCESS_VRDE_HOST, and
RDP_ACCESS_VRDE_PORT.
RDP_ACCESS_VRDE_PORT. In public mode, RDP_ACCESS_IPV6_BIND and
RDP_ACCESS_IPV6_FORWARD (auto, always, or never) control the reversible
IPv6-to-IPv4 listener for dual-stack hostnames. Set
RDP_ACCESS_GUACD_LOG_LEVEL=debug for temporary protocol diagnostics.
EOF
}
@@ -56,6 +65,16 @@ safe_name() {
case $2 in ''|*[!A-Za-z0-9.-]*) fail "$1 contains unsupported characters" ;; esac
}
safe_profile() {
case $2 in
[a-z0-9]* ) ;;
* ) fail "$1 must start with a lowercase alphanumeric" ;;
esac
case $2 in
*[!a-z0-9-]*) fail "$1 must contain only lowercase letters, digits, and hyphens" ;;
esac
}
safe_port() {
case $2 in ''|*[!0-9]*) fail "$1 must be a port number" ;; esac
[ "$2" -ge 1024 ] && [ "$2" -le 65535 ] || fail "$1 must be between 1024 and 65535"
@@ -69,6 +88,7 @@ compose() {
RDP_ACCESS_RUNTIME_DIR=$runtime_dir \
RDP_ACCESS_BIND=$RDP_ACCESS_BIND \
RDP_ACCESS_HTTP_PORT=$RDP_ACCESS_HTTP_PORT \
RDP_ACCESS_GUACD_LOG_LEVEL=$RDP_ACCESS_GUACD_LOG_LEVEL \
RDP_ACCESS_CERT_NAME=$RDP_ACCESS_PUBLIC_HOST \
RDP_ACCESS_CERT_SAN=$cert_san \
RDP_ACCESS_HOST_UID=$(id -u) \
@@ -76,15 +96,158 @@ compose() {
docker compose --project-name "$project" -f "$compose_file" "$@"
}
socket_path=$runtime_dir/ssh-control.socket
session_file=$runtime_dir/session.env
cert_name_file=$runtime_dir/cert-name
wait_for_gateway() {
command -v curl >/dev/null 2>&1 || fail "curl is required for the Guacamole readiness check"
attempts=0
while [ "$attempts" -lt 60 ]; do
# The listener is bound on the controller, even when the public
# address is 0.0.0.0. Ignore the self-signed certificate here: this
# check is only proving that nginx can reach the Guacamole servlet.
if curl -kfsS --connect-timeout 1 --max-time 3 \
"https://127.0.0.1:$RDP_ACCESS_HTTP_PORT/guacamole/" \
>/dev/null 2>&1; then
return 0
fi
attempts=$((attempts + 1))
sleep 1
done
return 1
}
stop_tunnel() {
mkdir -p "$runtime_dir"
: >"$tunnel_stop_file"
if [ -S "$socket_path" ]; then
ssh -S "$socket_path" -O exit "$RVBOX_TEST_VBOX_HOST" >/dev/null 2>&1 || true
ssh -S "$socket_path" -O exit -o ConnectTimeout=5 "$RVBOX_TEST_VBOX_HOST" >/dev/null 2>&1 || true
fi
if [ -f "$tunnel_pid_file" ]; then
tunnel_pid=$(cat "$tunnel_pid_file" 2>/dev/null || true)
case $tunnel_pid in
''|*[!0-9]*) ;;
*)
# Never trust a stale numeric PID file. Confirm that it is
# still this helper's watchdog before signalling it.
tunnel_cmd=$(ps -p "$tunnel_pid" -o args= 2>/dev/null || true)
case $tunnel_cmd in
*rdp-tunnel-watchdog*"$tunnel_stop_file"*)
kill "$tunnel_pid" >/dev/null 2>&1 || true
# The watchdog normally exits within its two-second
# retry interval. Do not leave a stale reconnect loop
# behind, but only force the verified process.
sleep 1
tunnel_cmd=$(ps -p "$tunnel_pid" -o args= 2>/dev/null || true)
case $tunnel_cmd in
*rdp-tunnel-watchdog*"$tunnel_stop_file"*) kill -KILL "$tunnel_pid" >/dev/null 2>&1 || true ;;
esac
;;
esac
;;
esac
fi
rm -f "$socket_path"
rm -f "$tunnel_pid_file" "$tunnel_stop_file"
}
ipv6_process_matches() {
ipv6_pid=$1
ipv6_cmd=$(ps -p "$ipv6_pid" -o args= 2>/dev/null || true)
case $ipv6_cmd in
*"TCP6-LISTEN:$RDP_ACCESS_HTTP_PORT"*"TCP4:127.0.0.1:$RDP_ACCESS_HTTP_PORT"*) return 0 ;;
*) return 1 ;;
esac
}
ipv6_listener_active() {
command -v ss >/dev/null 2>&1 || return 1
ss -H -ltn6 2>/dev/null | awk -v port=":$RDP_ACCESS_HTTP_PORT" \
'{ if ($4 ~ (port "$")) found=1 } END { exit !found }'
}
stop_ipv6_forward() {
[ -f "$ipv6_pid_file" ] || return 0
ipv6_pid=$(cat "$ipv6_pid_file" 2>/dev/null || true)
case $ipv6_pid in
''|*[!0-9]*) ;;
*)
if ipv6_process_matches "$ipv6_pid"; then
kill "$ipv6_pid" >/dev/null 2>&1 || true
sleep 1
if ipv6_process_matches "$ipv6_pid"; then
kill -KILL "$ipv6_pid" >/dev/null 2>&1 || true
fi
fi
;;
esac
rm -f "$ipv6_pid_file"
}
start_ipv6_forward() {
[ "$RDP_ACCESS_BIND" = 0.0.0.0 ] || return 0
case $RDP_ACCESS_IPV6_FORWARD in
never) stop_ipv6_forward; return 0 ;;
auto|always) ;;
*) fail "RDP_ACCESS_IPV6_FORWARD must be auto, always, or never" ;;
esac
if ! command -v socat >/dev/null 2>&1; then
if [ "$RDP_ACCESS_IPV6_FORWARD" = always ]; then
fail "socat is required for RDP_ACCESS_IPV6_FORWARD=always"
fi
printf 'rdp-access: warning: socat unavailable; IPv6 forwarding disabled (IPv4 remains active)\n' >&2
return 0
fi
if [ -f "$ipv6_pid_file" ] && ipv6_process_matches "$(cat "$ipv6_pid_file" 2>/dev/null || true)" && ipv6_listener_active; then
return 0
fi
if ipv6_listener_active; then
# Do not kill an unrelated listener. The operator can inspect it via
# status; down only stops a PID that this helper started and verified.
if [ "$RDP_ACCESS_IPV6_FORWARD" = always ]; then
fail "IPv6 port $RDP_ACCESS_HTTP_PORT is already occupied by an unowned listener"
fi
printf 'rdp-access: IPv6 port %s is already occupied; retaining the existing listener\n' \
"$RDP_ACCESS_HTTP_PORT" >&2
return 0
fi
stop_ipv6_forward
umask 077
mkdir -p "$runtime_dir"
: >"$ipv6_log_file"
chmod 600 "$ipv6_log_file"
setsid nohup socat \
"TCP6-LISTEN:$RDP_ACCESS_HTTP_PORT,ipv6only=1,reuseaddr,fork,bind=$RDP_ACCESS_IPV6_BIND" \
"TCP4:127.0.0.1:$RDP_ACCESS_HTTP_PORT" \
</dev/null >>"$ipv6_log_file" 2>&1 &
ipv6_pid=$!
printf '%s\n' "$ipv6_pid" >"$ipv6_pid_file"
chmod 600 "$ipv6_pid_file"
attempts=0
while [ "$attempts" -lt 10 ]; do
if ipv6_process_matches "$ipv6_pid" && ipv6_listener_active; then
return 0
fi
attempts=$((attempts + 1))
sleep 1
done
stop_ipv6_forward
if [ "$RDP_ACCESS_IPV6_FORWARD" = always ]; then
fail "could not start IPv6 forwarding listener; inspect $ipv6_log_file"
fi
printf 'rdp-access: warning: IPv6 forwarding could not start (IPv4 remains active); inspect %s\n' \
"$ipv6_log_file" >&2
}
ipv6_forward_status() {
if [ -f "$ipv6_pid_file" ] && ipv6_process_matches "$(cat "$ipv6_pid_file" 2>/dev/null || true)" && ipv6_listener_active; then
printf 'ipv6_forward=active\n'
printf 'ipv6_forward_supervisor=helper\n'
elif ipv6_listener_active; then
printf 'ipv6_forward=active_external\n'
printf 'ipv6_forward_supervisor=external\n'
elif [ "$RDP_ACCESS_BIND" != 0.0.0.0 ] || [ "$RDP_ACCESS_IPV6_FORWARD" = never ]; then
printf 'ipv6_forward=disabled\n'
else
printf 'ipv6_forward=inactive\n'
fi
}
gateway_for_network() {
@@ -101,7 +264,9 @@ password_hash_from_terminal() {
password=
restore_tty=false
if [ "$password_stdin" = true ]; then
IFS= read -r password || fail "could not read password from stdin"
# Password files commonly omit a final newline. POSIX read returns
# non-zero at that EOF even after assigning the final nonempty line.
IFS= read -r password || [ -n "$password" ] || fail "could not read password from stdin"
else
[ -t 0 ] || fail "stdin is not a terminal; use --web-password-stdin"
printf 'Fixture password for %s: ' "$RDP_ACCESS_WEB_USER" >&2
@@ -150,19 +315,111 @@ render_mapping() {
printf 'url=https://%s:%s/guacamole/\n' "$RDP_ACCESS_PUBLIC_HOST" "$RDP_ACCESS_HTTP_PORT" >"$session_file"
printf 'docker_gateway=%s\n' "$docker_gateway" >>"$session_file"
printf 'tunnel_port=%s\n' "$RDP_ACCESS_TUNNEL_PORT" >>"$session_file"
printf 'bind=%s\n' "$RDP_ACCESS_BIND" >>"$session_file"
printf 'http_port=%s\n' "$RDP_ACCESS_HTTP_PORT" >>"$session_file"
printf 'public_host=%s\n' "$RDP_ACCESS_PUBLIC_HOST" >>"$session_file"
printf 'web_user=%s\n' "$RDP_ACCESS_WEB_USER" >>"$session_file"
printf 'rdp_user=%s\n' "$RDP_ACCESS_RDP_USER" >>"$session_file"
chmod 600 "$session_file"
}
session_has() {
grep -Fqx "$1=$2" "$session_file"
}
session_matches_current() {
[ -f "$session_file" ] || return 1
session_has url "https://$RDP_ACCESS_PUBLIC_HOST:$RDP_ACCESS_HTTP_PORT/guacamole/" || return 1
session_has docker_gateway "$docker_gateway" || return 1
session_has tunnel_port "$RDP_ACCESS_TUNNEL_PORT" || return 1
session_has bind "$RDP_ACCESS_BIND" || return 1
session_has http_port "$RDP_ACCESS_HTTP_PORT" || return 1
session_has public_host "$RDP_ACCESS_PUBLIC_HOST" || return 1
session_has web_user "$RDP_ACCESS_WEB_USER" || return 1
session_has rdp_user "$RDP_ACCESS_RDP_USER"
}
start_tunnel() {
stop_tunnel
ssh -M -S "$socket_path" -fN \
-o BatchMode=yes \
-o ExitOnForwardFailure=yes \
-o ServerAliveInterval=30 \
-o ServerAliveCountMax=3 \
-L "$docker_gateway:$RDP_ACCESS_TUNNEL_PORT:$RDP_ACCESS_VRDE_HOST:$RDP_ACCESS_VRDE_PORT" \
"$RVBOX_TEST_VBOX_HOST"
ssh -S "$socket_path" -O check "$RVBOX_TEST_VBOX_HOST" >/dev/null 2>&1 || fail "private VRDE tunnel did not start"
umask 077
: >"$tunnel_log_file"
chmod 600 "$tunnel_log_file"
# Keep the SSH master outside the short-lived controller process. Helium
# may drop an idle or metered route; the watchdog reconnects it and keeps
# the same local listener for Guacamole. The stop marker is deliberately
# inside .runtime so down/clean can terminate the loop deterministically.
watchdog_script='
set -eu
stop_file=$1
socket=$2
host=$3
bind=$4
tunnel_port=$5
vrde_host=$6
vrde_port=$7
while [ ! -e "$stop_file" ]; do
if [ -S "$socket" ] && ssh -S "$socket" -O check "$host" >/dev/null 2>&1; then
sleep 2
continue
fi
rm -f "$socket"
ssh -M -S "$socket" -fN \
-o BatchMode=yes \
-o ExitOnForwardFailure=yes \
-o ServerAliveInterval=30 \
-o ServerAliveCountMax=3 \
-L "$bind:$tunnel_port:$vrde_host:$vrde_port" \
"$host" || true
sleep 2
done
'
if command -v setsid >/dev/null 2>&1; then
setsid nohup sh -c "$watchdog_script" rdp-tunnel-watchdog \
"$tunnel_stop_file" "$socket_path" "$RVBOX_TEST_VBOX_HOST" \
"$docker_gateway" "$RDP_ACCESS_TUNNEL_PORT" "$RDP_ACCESS_VRDE_HOST" \
"$RDP_ACCESS_VRDE_PORT" </dev/null >>"$tunnel_log_file" 2>&1 &
else
nohup sh -c "$watchdog_script" rdp-tunnel-watchdog \
"$tunnel_stop_file" "$socket_path" "$RVBOX_TEST_VBOX_HOST" \
"$docker_gateway" "$RDP_ACCESS_TUNNEL_PORT" "$RDP_ACCESS_VRDE_HOST" \
"$RDP_ACCESS_VRDE_PORT" </dev/null >>"$tunnel_log_file" 2>&1 &
fi
tunnel_watchdog_pid=$!
printf '%s\n' "$tunnel_watchdog_pid" >"$tunnel_pid_file"
chmod 600 "$tunnel_pid_file"
# A listener is useful only after the SSH master has authenticated and
# bound the Docker-network gateway. Give the reconnect loop a bounded
# window, then fail with a clear recovery path.
attempts=0
while [ "$attempts" -lt 20 ]; do
if [ -S "$socket_path" ] && ssh -S "$socket_path" -O check "$RVBOX_TEST_VBOX_HOST" >/dev/null 2>&1; then
return 0
fi
if command -v nc >/dev/null 2>&1 && nc -z -w 2 "$docker_gateway" "$RDP_ACCESS_TUNNEL_PORT" >/dev/null 2>&1; then
return 0
fi
attempts=$((attempts + 1))
sleep 1
done
stop_tunnel
return 1
}
tunnel_listener_active() {
[ -n "${docker_gateway-}" ] || return 1
command -v nc >/dev/null 2>&1 || return 1
nc -z -w 2 "$docker_gateway" "$RDP_ACCESS_TUNNEL_PORT" >/dev/null 2>&1
}
tunnel_control_active() {
[ -S "$socket_path" ] || return 1
ssh -S "$socket_path" -O check "$RVBOX_TEST_VBOX_HOST" >/dev/null 2>&1
}
tunnel_active() {
tunnel_control_active || tunnel_listener_active
}
action=${1-}
@@ -175,6 +432,7 @@ password_stdin=false
remove_images=false
while [ "$#" -gt 0 ]; do
case $1 in
--profile) [ "$#" -ge 2 ] || fail "--profile needs a value"; RDP_ACCESS_PROFILE=$2; shift 2 ;;
--bind) [ "$#" -ge 2 ] || fail "--bind needs a value"; RDP_ACCESS_BIND=$2; shift 2 ;;
--http-port) [ "$#" -ge 2 ] || fail "--http-port needs a value"; RDP_ACCESS_HTTP_PORT=$2; shift 2 ;;
--tunnel-port) [ "$#" -ge 2 ] || fail "--tunnel-port needs a value"; RDP_ACCESS_TUNNEL_PORT=$2; shift 2 ;;
@@ -188,6 +446,26 @@ while [ "$#" -gt 0 ]; do
esac
done
safe_profile RDP_ACCESS_PROFILE "$RDP_ACCESS_PROFILE"
profile=$RDP_ACCESS_PROFILE
if [ "$profile" = default ]; then
runtime_dir=$runtime_root
project=rvbox-rdp-access
else
runtime_dir=$runtime_root/$profile
project=rvbox-rdp-access-$profile
fi
[ ! -L "$runtime_root" ] || fail "runtime root must not be a symlink"
[ ! -L "$runtime_dir" ] || fail "runtime directory must not be a symlink"
socket_path=$runtime_dir/ssh-control.socket
session_file=$runtime_dir/session.env
cert_name_file=$runtime_dir/cert-name
tunnel_stop_file=$runtime_dir/tunnel.stop
tunnel_pid_file=$runtime_dir/tunnel-watchdog.pid
tunnel_log_file=$runtime_dir/tunnel.log
ipv6_pid_file=$runtime_dir/ipv6-forward.pid
ipv6_log_file=$runtime_dir/ipv6-forward.log
safe_bind "$RDP_ACCESS_BIND"
safe_port RDP_ACCESS_HTTP_PORT "$RDP_ACCESS_HTTP_PORT"
safe_port RDP_ACCESS_TUNNEL_PORT "$RDP_ACCESS_TUNNEL_PORT"
@@ -198,6 +476,8 @@ safe_name RDP_ACCESS_RDP_USER "$RDP_ACCESS_RDP_USER"
safe_name RVBOX_TEST_VBOX_HOST "$RVBOX_TEST_VBOX_HOST"
case $RDP_ACCESS_VRDE_HOST in 127.0.0.1|localhost) ;; *) fail "RDP_ACCESS_VRDE_HOST must be 127.0.0.1 or localhost" ;; esac
safe_port RDP_ACCESS_VRDE_PORT "$RDP_ACCESS_VRDE_PORT"
case $RDP_ACCESS_IPV6_BIND in ::|::1) ;; *) fail "RDP_ACCESS_IPV6_BIND must be :: or ::1" ;; esac
case $RDP_ACCESS_IPV6_FORWARD in auto|always|never) ;; *) fail "RDP_ACCESS_IPV6_FORWARD must be auto, always, or never" ;; esac
case $RDP_ACCESS_PUBLIC_HOST in
*[!0-9.]* ) cert_san="DNS:$RDP_ACCESS_PUBLIC_HOST" ;;
@@ -218,7 +498,24 @@ case $action in
docker compose version >/dev/null
assert_fixture_running
if compose ps -q | grep -q .; then
fail "gateway already exists; use status or down first"
# A short-lived controller (or a dropped SSH route) can leave the
# Compose stack running after its tunnel has disappeared. Reuse
# that stack and repair only the private forwarding path instead
# of forcing the operator to tear down a usable Guacamole session.
mkdir -p "$runtime_dir"
docker_gateway=$(gateway_for_network) || fail "could not determine private Docker gateway"
session_matches_current || fail "gateway already exists with a different configuration; run down first, then run up with the desired options"
if tunnel_active; then
printf 'Guacamole stack and private VRDE tunnel are already active.\n'
elif start_tunnel; then
printf 'Existing Guacamole stack reused; private VRDE tunnel restored.\n'
else
fail "existing Guacamole stack was found but the private VRDE tunnel could not be restored; inspect $tunnel_log_file"
fi
wait_for_gateway || fail "Guacamole stack is running but its HTTPS endpoint did not become ready; inspect Compose logs"
start_ipv6_forward
if [ -f "$session_file" ]; then sed -n '1p' "$session_file"; fi
exit 0
fi
mkdir -p "$runtime_dir"
compose up -d guacd
@@ -238,19 +535,43 @@ case $action in
compose down --remove-orphans
fail "could not start Guacamole gateway"
fi
if ! wait_for_gateway; then
fail "Guacamole gateway started but its HTTPS endpoint did not become ready; inspect Compose logs"
fi
start_ipv6_forward
printf 'Guacamole is ready at https://%s:%s/guacamole/\n' "$RDP_ACCESS_PUBLIC_HOST" "$RDP_ACCESS_HTTP_PORT"
printf 'Accept the self-signed certificate warning, then sign in as %s with the fixture password.\n' "$RDP_ACCESS_WEB_USER"
;;
status)
[ "$#" -eq 0 ] || { usage >&2; fail "status accepts no options"; }
printf 'profile=%s project=%s runtime=%s\n' "$profile" "$project" "$runtime_dir"
"$repo_root/scripts/windows/test-host" status || true
if [ -f "$session_file" ]; then sed -n '1p' "$session_file"; fi
compose ps
if [ -S "$socket_path" ] && ssh -S "$socket_path" -O check "$RVBOX_TEST_VBOX_HOST" >/dev/null 2>&1; then
docker_gateway=$(gateway_for_network 2>/dev/null || true)
if tunnel_control_active; then
printf 'private_tunnel=active\n'
printf 'private_tunnel_supervisor=helper\n'
elif tunnel_listener_active; then
# A tunnel started by an interactive shell or another supervisor
# has no control socket owned by this helper, but it is still a
# valid path when the Docker listener is reachable.
printf 'private_tunnel=active_external\n'
printf 'private_tunnel_supervisor=external\n'
else
printf 'private_tunnel=inactive\n'
fi
ipv6_forward_status
;;
repair)
[ "$#" -eq 0 ] || { usage >&2; fail "repair accepts no options"; }
docker version >/dev/null
docker compose version >/dev/null
assert_fixture_running
[ -n "$(compose ps -q guacd 2>/dev/null || true)" ] || fail "Guacamole stack is not running; run up first"
compose restart guacd >/dev/null
wait_for_gateway || fail "Guacamole gateway did not become ready after guacd repair"
printf 'guacd restarted; stale VRDE workers released. Close any old browser tab and sign in again.\n'
;;
url)
[ "$#" -eq 0 ] || { usage >&2; fail "url accepts no options"; }
@@ -262,15 +583,20 @@ case $action in
;;
down)
[ "$#" -eq 0 ] || { usage >&2; fail "down accepts no options"; }
stop_ipv6_forward
stop_tunnel
compose down --remove-orphans || true
printf 'Temporary gateway and private tunnel stopped; generated certificate and password hash retained in %s.\n' "$runtime_dir"
;;
clean)
[ "$#" -eq 0 ] || { usage >&2; fail "clean accepts only --images"; }
stop_ipv6_forward
stop_tunnel
compose down --remove-orphans || true
case $runtime_dir in "$helper_dir"/.runtime) rm -rf "$runtime_dir" ;; *) fail "unsafe runtime path" ;; esac
case "$runtime_dir" in
"$runtime_root"|"$runtime_root/$profile") rm -rf "$runtime_dir" ;;
*) fail "unsafe runtime path" ;;
esac
if [ "$remove_images" = true ]; then
docker image rm guacamole/guacamole:1.6.0 guacamole/guacd:1.6.0 nginx:1.27-alpine >/dev/null 2>&1 || true
fi
+1
View File
@@ -15,6 +15,7 @@
<param name="enable-wallpaper">false</param>
<param name="enable-theming">false</param>
<param name="enable-font-smoothing">false</param>
<param name="color-depth">24</param>
<param name="disable-gfx">true</param>
</connection>
</authorize>
+12
View File
@@ -49,6 +49,18 @@ func TestNativeFixtureAssets_HP_HARNESS_20(t *testing.T) {
t.Fatalf("native test-host is missing compressed transfer contract %q", required)
}
}
rdpAccess := read("test/rdp-access/rdp-access")
for _, required := range []string{"RDP_ACCESS_PROFILE", "--profile", "rvbox-rdp-access-$profile", "runtime_root", "ipv6_forward_status", "compose restart guacd", "compose down --remove-orphans"} {
if !strings.Contains(rdpAccess, required) {
t.Fatalf("RDP helper is missing isolated lifecycle contract %q", required)
}
}
rdpReadme := read("test/rdp-access/README.md")
for _, required := range []string{"Current fixture quick start", "RDP_ACCESS_PROFILE=vm-b", "Partial stack or changed endpoint", "Inspect, repair, stop, or clean that instance"} {
if !strings.Contains(rdpReadme, required) {
t.Fatalf("RDP helper documentation is missing %q", required)
}
}
compose := read("test/linux-server/compose.yaml")
for _, required := range []string{"../../bin/rvbox-server", "nginx:1.27-alpine", "rvbox.native.run_id", "RVBOX_NATIVE_RUNTIME_DIR", "certgen", "networks: [native]"} {
if !strings.Contains(compose, required) {
+89
View File
@@ -0,0 +1,89 @@
// Command releaseicon creates the small deterministic PNG used only by the
// Windows release resource generator. Keeping the source rather than a binary
// asset makes the icon reproducible in the pinned toolchain.
package main
import (
"flag"
"image"
"image/color"
"image/png"
"os"
"path/filepath"
)
const canvasSize = 256
func main() {
output := flag.String("out", "", "absolute PNG output path")
flag.Parse()
if *output == "" || !filepath.IsAbs(*output) {
panic("--out must be an absolute path")
}
if err := os.MkdirAll(filepath.Dir(*output), 0o700); err != nil {
panic(err)
}
file, err := os.OpenFile(*output, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0o600)
if err != nil {
panic(err)
}
defer file.Close()
image := image.NewNRGBA(image.Rect(0, 0, canvasSize, canvasSize))
fill(image, color.NRGBA{R: 15, G: 28, B: 48, A: 255})
fillRoundedRect(image, 18, 18, 238, 238, 42, color.NRGBA{R: 28, G: 56, B: 92, A: 255})
// A compact outlined box and terminal prompt remain recognizable at the
// 16-pixel size Windows Explorer generates from this source image.
strokeRect(image, 57, 62, 199, 194, 12, color.NRGBA{R: 54, G: 211, B: 189, A: 255})
fillRect(image, 79, 105, 96, 121, color.NRGBA{R: 240, G: 249, B: 255, A: 255})
fillRect(image, 95, 121, 112, 137, color.NRGBA{R: 240, G: 249, B: 255, A: 255})
fillRect(image, 79, 137, 96, 153, color.NRGBA{R: 240, G: 249, B: 255, A: 255})
fillRect(image, 126, 150, 176, 166, color.NRGBA{R: 240, G: 249, B: 255, A: 255})
if err := png.Encode(file, image); err != nil {
panic(err)
}
}
func fill(image *image.NRGBA, value color.NRGBA) {
for y := 0; y < canvasSize; y++ {
for x := 0; x < canvasSize; x++ {
image.SetNRGBA(x, y, value)
}
}
}
func fillRoundedRect(image *image.NRGBA, left, top, right, bottom, radius int, value color.NRGBA) {
for y := top; y < bottom; y++ {
for x := left; x < right; x++ {
dx, dy := 0, 0
if x < left+radius {
dx = left + radius - x
} else if x >= right-radius {
dx = x - (right - radius - 1)
}
if y < top+radius {
dy = top + radius - y
} else if y >= bottom-radius {
dy = y - (bottom - radius - 1)
}
if dx*dx+dy*dy <= radius*radius {
image.SetNRGBA(x, y, value)
}
}
}
}
func strokeRect(image *image.NRGBA, left, top, right, bottom, width int, value color.NRGBA) {
fillRect(image, left, top, right, top+width, value)
fillRect(image, left, bottom-width, right, bottom, value)
fillRect(image, left, top, left+width, bottom, value)
fillRect(image, right-width, top, right, bottom, value)
}
func fillRect(image *image.NRGBA, left, top, right, bottom int, value color.NRGBA) {
for y := top; y < bottom; y++ {
for x := left; x < right; x++ {
image.SetNRGBA(x, y, value)
}
}
}