Compare commits

...
6 Commits
30 changed files with 882 additions and 83 deletions
+11 -4
View File
@@ -6,6 +6,7 @@ import (
"errors"
"flag"
"fmt"
"io"
"log"
"net"
"net/http"
@@ -58,6 +59,12 @@ func run(configPath string) error {
if err != nil {
return err
}
runtimeLog, err := observability.OpenRotatingFile(configured.Observability.LogFile, configured.Observability.LogMaxBytes, configured.Observability.LogMaxFiles)
if err != nil {
return fmt.Errorf("open server rotating log: %w", err)
}
defer runtimeLog.Close()
logger := log.New(io.MultiWriter(os.Stderr, observability.FormatLog(runtimeLog, configured.Observability.LogFormat)), "rvbox-server: ", log.LstdFlags|log.LUTC)
persistence, err := store.Open(context.Background(), store.Options{
DataDir: configured.Server.DataDir, BusyTimeout: configured.Storage.SQLiteBusyTimeout,
SegmentTargetSize: configured.Storage.SegmentTargetBytes,
@@ -84,7 +91,7 @@ func run(configPath string) error {
if configured.Observability.Listen != "" {
healthListener, err = net.Listen("tcp", configured.Observability.Listen)
if err != nil {
log.Printf("rvbox-server: observability endpoint unavailable (continuing without it): %v", err)
logger.Printf("observability endpoint unavailable (continuing without it): %v", err)
}
}
if healthListener != nil {
@@ -104,10 +111,10 @@ func run(configPath string) error {
Scope: store.IncidentScopeGlobal, ScopeKey: "server-startup-recovery", Summary: "startup storage recovery failed",
Evidence: []byte(recoverErr.Error()), AutomaticallyRepairable: false,
}); incidentErr != nil {
log.Printf("rvbox-server: could not persist recovery incident: %v", incidentErr)
logger.Printf("could not persist recovery incident: %v", incidentErr)
}
}
log.Printf("rvbox-server: storage recovery left readiness disabled: %v", recoverErr)
logger.Printf("storage recovery left readiness disabled: %v", recoverErr)
return
}
health.SetReady(true)
@@ -143,7 +150,7 @@ func run(configPath string) error {
}
defer rpcListener.Close()
if configured.JSONRPC.NonLoopbackBind {
log.Printf("rvbox-server: WARNING JSON-RPC is unauthenticated and bound to non-loopback address %s", configured.JSONRPC.Listen)
logger.Printf("WARNING JSON-RPC is unauthenticated and bound to non-loopback address %s", configured.JSONRPC.Listen)
}
rpcServer = &http.Server{Handler: control.NewJSONRPCHandler(controlService, int64(configured.Protocol.MaxJSONRPCBodyBytes)), ReadHeaderTimeout: configured.Flow.WriteDeadline}
}
+16
View File
@@ -150,6 +150,17 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ
if err != nil {
return err
}
runtimeLog, err := observability.OpenRotatingFile(configured.Observability.LogFile, configured.Observability.LogMaxBytes, configured.Observability.LogMaxFiles)
if err != nil {
return fmt.Errorf("open client rotating log: %w", err)
}
defer runtimeLog.Close()
formattedRuntimeLog := observability.FormatLog(runtimeLog, configured.Observability.LogFormat)
if diagnostics == nil {
diagnostics = formattedRuntimeLog
} else {
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 {
@@ -216,6 +227,11 @@ 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,
OnSessionError: func(sessionErr error) {
if diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox client session retry: %v\n", sessionErr)
}
},
}); runErr != nil && ctx.Err() == nil && diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox client session stopped: %v\n", runErr)
}
+23
View File
@@ -14,6 +14,7 @@ import (
clientwindows "github.com/rvbox/rvbox/internal/client/supervisor/windows"
"github.com/rvbox/rvbox/internal/client/windowsservice"
"github.com/rvbox/rvbox/internal/client/windowstray"
"golang.org/x/sys/windows"
"golang.org/x/sys/windows/svc"
)
@@ -90,6 +91,16 @@ func openServiceDiagnostics(configPath string, diagnostics io.Writer) (io.Writer
}
return diagnostics, func() {}
}
// Service diagnostics are intentionally not private spool data. Preserve a
// protected SYSTEM/Administrators DACL so an operator can diagnose an SCM
// startup or reconnect failure without weakening access to command state.
if err := applyServiceDiagnosticsACL(path); err != nil {
_ = file.Close()
if diagnostics == nil {
return io.Discard, func() {}
}
return diagnostics, func() {}
}
if diagnostics == nil {
return file, func() { _ = file.Close() }
}
@@ -99,6 +110,18 @@ func openServiceDiagnostics(configPath string, diagnostics io.Writer) (io.Writer
return io.MultiWriter(file, diagnostics), func() { _ = file.Close() }
}
func applyServiceDiagnosticsACL(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
}
return windows.SetNamedSecurityInfo(path, windows.SE_FILE_OBJECT, windows.DACL_SECURITY_INFORMATION|windows.PROTECTED_DACL_SECURITY_INFORMATION, nil, nil, dacl, nil)
}
func runTray(configPath string, diagnostics io.Writer) error {
_ = configPath // the tray obtains the canonical paths from the service.
return windowstray.Run(context.Background(), diagnostics)
+3
View File
@@ -27,6 +27,9 @@ services:
server:
image: "${RVBOX_SERVER_IMAGE:?set RVBOX_SERVER_IMAGE to a pinned rvbox-server image}"
# The init container is the only root process. Its exact 65532 ownership
# setup lets the long-lived server run without host/container root power.
user: "65532:65532"
restart: unless-stopped
command: ["--config", "/etc/rvbox/server.toml"]
read_only: true
+6
View File
@@ -19,3 +19,9 @@ enabled = false
[observability]
# nginx exposes only /livez and /readyz, not the metrics listener itself.
listen = "0.0.0.0:6901"
# The server retains its own bounded diagnostics beside durable state.
log_file = "/var/lib/rvbox-server/logs/rvbox-server.log"
# Rotate the current log after this many bytes (10 MiB).
log_max_bytes = 10485760
# Number of sealed rotated log files to retain.
log_max_files = 5
+5 -2
View File
@@ -97,8 +97,11 @@ Scheduler is not used.
4. The server queues work while a client is offline and dispatches it only when
the active session advertises capacity. A replacement session immediately
performs bidirectional reconciliation between the server's non-terminal set
and the client's complete retained-command set. The server returns explicit
local terminate/discard decisions before new dispatch begins.
and the client's complete retained-command set. That comparison uses one
server receipt-time cutoff captured before session registration, so commands
admitted during the handshake remain fresh queued work rather than false
missing-client contradictions. The server returns explicit local
terminate/discard decisions before new dispatch begins.
## Command model
+6
View File
@@ -1751,6 +1751,12 @@ a lost final `EventAck` followed by server retention. Write this idempotent
result before dispatch. A client with unresolved essential-store corruption
cannot complete reconciliation or accept work.
Capture one server receipt-time cutoff immediately before registering the live
session, and use that identical cutoff when building `ReconcileRequest` and
applying its `ReconcileSnapshot`. A command admitted after that cutoff is new
work: it must remain queued for post-reconciliation dispatch, never be treated
as absent client evidence merely because its admission raced the handshake.
Implement reconciliation as this explicit matrix:
| Server state | Client evidence | Durable result |
+27 -5
View File
@@ -15,6 +15,18 @@ Run focused unit tests with:
scripts/test-unit --package ./internal/domain --run UUIDv7 --race
```
Run one bounded hostile-input fuzz target in the same pinned container with:
```sh
scripts/test-fuzz --package ./internal/agentproto --name FuzzDecodeEnvelopeBounded_SEC_PROTO_01 --time 30s
scripts/test-fuzz --package ./internal/server/control --name FuzzDecodeJSONRPCRequestBounded_SEC_CTL_01 --time 30s
```
The fuzz command disables ordinary tests for that invocation and runs exactly
one target. A crash stores Go's minimized reproducer in the affected package's
fuzz corpus, so the normal test gate exercises it as a deterministic seed
after it is reviewed and committed.
Build all supported binaries without installing `make` or Go on the host:
```sh
@@ -70,8 +82,12 @@ default). Helium hosts the VM only; it does not host any RVBox server
containers. The self-signed server certificate is intentionally accepted by
the v1 client without a test CA. It then drives the installed SCM service
through the server's real Unix control socket and verifies every Windows
execution context. The tagged binary's controlled pre-launch failures are
limited to the test fixture; a release binary rejects that switch.
execution context, ordered stdin close, TERM delivery, and a server-process
restart while a command is running. The restart check preserves the server
state volume, waits for a new reconciled WSS session, then proves that the same
command can receive its terminal signal; it covers reconnect without treating
the old session as valid. The tagged binary's controlled pre-launch failures
are limited to the test fixture; a release binary rejects that switch.
Successful runs collect bounded artifacts, remove only their labeled Compose
project, and restore the exact clean snapshot. A failed or --keep run stays
@@ -222,7 +238,7 @@ that mode-600 file; the controller never puts it on a command line, manifest,
log, or artifact.
The native lifecycle is `status`, `prepare`, `stage`, `install`, `run`,
`collect`, `stop`, and `reset`. `prepare` verifies the VM and snapshot UUIDs,
`collect`, `logs`, `stop`, and `reset`. `prepare` verifies the VM and snapshot UUIDs,
restores the clean baseline, starts headless, waits for Guest Additions, and
proves that `RVBoxClient` is absent. `stage` copies a versioned non-secret test
bundle through a run-specific host directory to a run-specific guest directory.
@@ -230,8 +246,14 @@ bundle through a run-specific host directory to a run-specific guest directory.
the real `rvbox.exe --install-service` path and proves completion through SCM.
`run` is for reconfiguration/restart scenarios after that first installation.
Neither action invokes the GUI-subsystem executable directly with the normal
Guest Control account. `collect` obtains only bounded/redacted artifacts, and
`reset` restores the exact clean baseline and leaves the VM powered off.
Guest Control account. `collect` obtains only bounded/redacted artifacts.
`logs` is the narrow read-only service-startup/client-diagnostics action for a
retained prepared run. The service diagnostic file grants access to SYSTEM and
local Administrators only; it contains no command spool data. `reset` restores
the exact clean baseline and leaves the VM powered off. It
first permits a bounded ACPI shutdown; if that hangs, it force-powers off only
the exact leased disposable VM before snapshot restoration. That intentional
state loss is confined to the test isolation boundary.
This service-driven protocol is required because VirtualBox Guest Control
7.2.16 does not reliably complete a direct GUI-subsystem `rvbox.exe` run;
+29
View File
@@ -42,6 +42,35 @@ func TestDecodeEnvelopeBoundaries_HP_PROTO_01(t *testing.T) {
}
}
func FuzzDecodeEnvelopeBounded_SEC_PROTO_01(f *testing.F) {
valid, err := proto.Marshal(&rvboxv1.AgentEnvelope{Payload: &rvboxv1.AgentEnvelope_ClientHello{ClientHello: validHello()}})
if err != nil {
f.Fatal(err)
}
f.Add(valid)
f.Add([]byte{0xff})
f.Fuzz(func(t *testing.T, data []byte) {
limits := DefaultLimits()
if len(data) > int(limits.MaxEnvelopeBytes)+1 {
return
}
_, _ = DecodeEnvelope(data, limits, rvboxv1.Platform_PLATFORM_WINDOWS)
})
}
func FuzzDecodeOutputChunkBounded_SEC_PROTO_02(f *testing.F) {
f.Add([]byte("plain"), uint8(rvboxv1.Compression_COMPRESSION_NONE), uint64(5))
f.Add([]byte{0x28, 0xb5, 0x2f, 0xfd}, uint8(rvboxv1.Compression_COMPRESSION_ZSTD), uint64(1))
f.Fuzz(func(t *testing.T, data []byte, compression uint8, rawBytes uint64) {
const limit = uint64(64 << 10)
if len(data) > int(limit) {
return
}
chunk := &rvboxv1.OutputChunk{Stream: rvboxv1.StreamKind_STREAM_STDOUT, Compression: rvboxv1.Compression(compression % 3), CompressedSize: uint64(len(data)), UncompressedSize: rawBytes % (limit + 2), Data: data}
_, _ = DecodeOutputChunk(chunk, limit)
})
}
func TestEnvelopeSessionFencingShape_BH_SES_01(t *testing.T) {
t.Parallel()
+29 -22
View File
@@ -35,6 +35,10 @@ type RunnerOptions struct {
OnSignal func(context.Context, Session, *rvboxv1.SignalCommand) error
OnScriptReady func(context.Context, Session, domain.UUID) error
OnTerminate func(context.Context, domain.UUID) error
// OnSessionError observes one failed dial, handshake, protocol, or active
// transport session before normal reconnect backoff. It must not block; the
// durable spool and retry policy remain owned by Run.
OnSessionError func(error)
// 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.
@@ -95,10 +99,13 @@ func Run(ctx context.Context, options RunnerOptions) error {
return nil
}
sessionStarted := options.Now()
_ = runOnce(ctx, options)
sessionErr := runOnce(ctx, options)
if ctx.Err() != nil {
return nil
}
if sessionErr != nil && options.OnSessionError != nil {
options.OnSessionError(sessionErr)
}
if options.Now().Sub(sessionStarted) >= options.Backoff.StableReset {
failures = 0
}
@@ -152,7 +159,7 @@ func RunOnce(ctx context.Context, options RunnerOptions) error {
func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
transport, err := options.Dial(ctx)
if err != nil {
return err
return fmt.Errorf("dial agent server: %w", err)
}
defer func() {
if closeErr := transport.Close(); resultErr == nil && closeErr != nil {
@@ -166,7 +173,7 @@ func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
}
session, err := Handshake(ctx, transport, options.Hello, limits)
if err != nil {
return err
return fmt.Errorf("agent handshake: %w", err)
}
snapshot, err := options.Store.ReconcileSnapshot(ctx)
if err != nil {
@@ -174,7 +181,7 @@ func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
}
result, err := Reconcile(ctx, transport, session, snapshot, limits)
if err != nil {
return err
return fmt.Errorf("exchange reconciliation: %w", err)
}
terminated, err := ApplyReconcileResult(ctx, options.Store, result, options.Now())
if err != nil {
@@ -183,16 +190,16 @@ func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
for _, issue := range terminated {
if options.OnTerminate != nil {
if err := options.OnTerminate(ctx, issue); err != nil {
return err
return fmt.Errorf("terminate reconciled command %s: %w", issue, err)
}
}
}
sentEvents := make(map[domain.UUID]uint64)
if err := replayEvents(ctx, transport, options.Store, session, snapshot, limits, sentEvents); err != nil {
return err
return fmt.Errorf("replay client events: %w", err)
}
if err := sendCapacity(ctx, transport, session, options.Hello, snapshot, limits); err != nil {
return err
return fmt.Errorf("advertise client capacity: %w", err)
}
return serveActive(ctx, transport, options, session, sentEvents)
}
@@ -323,20 +330,20 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
case <-ctx.Done():
return nil
case err := <-readErrors:
return err
return fmt.Errorf("read active agent frame: %w", err)
case issue := <-options.EventReady:
if issue == (domain.UUID{}) {
continue
}
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != nil {
return err
return fmt.Errorf("flush ready command %s events: %w", issue, err)
}
continue
case encoded = <-frames:
}
envelope, err := agentproto.DecodeEnvelope(encoded, limits, rvboxv1.Platform_PLATFORM_WINDOWS)
if err != nil {
return err
return fmt.Errorf("decode active agent frame: %w", err)
}
if envelope.GetSessionId() != session.ID || envelope.GetSessionGeneration() != session.Generation {
return ErrProtocolHandshake
@@ -344,25 +351,25 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
switch {
case envelope.GetCommandDispatch() != nil:
if err := handleDispatch(ctx, transport, options, session, envelope.GetCommandDispatch(), limits); err != nil {
return err
return fmt.Errorf("handle command dispatch %s: %w", envelope.GetCommandDispatch().GetIssueUuid(), err)
}
issue, parseErr := domain.ParseUUIDv7(envelope.GetCommandDispatch().GetIssueUuid())
if parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != nil {
return err
return fmt.Errorf("flush dispatched command %s events: %w", issue, err)
}
}
case envelope.GetEventAck() != nil:
if err := ApplyEventAck(ctx, options.Store, envelope.GetEventAck()); err != nil {
return err
return fmt.Errorf("apply event acknowledgement for %s: %w", envelope.GetEventAck().GetIssueUuid(), err)
}
case envelope.GetScriptChunk() != nil:
if err := handleScriptChunk(ctx, transport, options.Store, session, envelope.GetScriptChunk(), limits, options.Now, sent); err != nil {
return err
return fmt.Errorf("apply script chunk for %s: %w", envelope.GetScriptChunk().GetIssueUuid(), err)
}
case envelope.GetScriptCommit() != nil:
if err := handleScriptCommit(ctx, transport, options.Store, session, envelope.GetScriptCommit(), limits, options.Now, sent); err != nil {
return err
return fmt.Errorf("commit script for %s: %w", envelope.GetScriptCommit().GetIssueUuid(), err)
}
if options.OnScriptReady != nil {
issue, err := domain.ParseUUIDv7(envelope.GetScriptCommit().GetIssueUuid())
@@ -370,40 +377,40 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
return err
}
if err := options.OnScriptReady(ctx, session, issue); err != nil {
return err
return fmt.Errorf("start committed script %s: %w", issue, err)
}
}
case envelope.GetStdinWrite() != nil:
if options.OnStdin != nil {
if err := options.OnStdin(ctx, session, envelope.GetStdinWrite()); err != nil {
return err
return fmt.Errorf("append stdin for %s: %w", envelope.GetStdinWrite().GetIssueUuid(), err)
}
}
if issue, parseErr := domain.ParseUUIDv7(envelope.GetStdinWrite().GetIssueUuid()); parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != nil {
return err
return fmt.Errorf("flush stdin command %s events: %w", issue, err)
}
}
case envelope.GetCloseStdin() != nil:
if options.OnCloseStdin != nil {
if err := options.OnCloseStdin(ctx, session, envelope.GetCloseStdin()); err != nil {
return err
return fmt.Errorf("close stdin for %s: %w", envelope.GetCloseStdin().GetIssueUuid(), err)
}
}
if issue, parseErr := domain.ParseUUIDv7(envelope.GetCloseStdin().GetIssueUuid()); parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != nil {
return err
return fmt.Errorf("flush closed-stdin command %s events: %w", issue, err)
}
}
case envelope.GetSignalCommand() != nil:
if options.OnSignal != nil {
if err := options.OnSignal(ctx, session, envelope.GetSignalCommand()); err != nil {
return err
return fmt.Errorf("signal command %s: %w", envelope.GetSignalCommand().GetIssueUuid(), err)
}
}
if issue, parseErr := domain.ParseUUIDv7(envelope.GetSignalCommand().GetIssueUuid()); parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != nil {
return err
return fmt.Errorf("flush signalled command %s events: %w", issue, err)
}
}
case envelope.GetError() != nil:
+30
View File
@@ -72,6 +72,36 @@ func TestRunnerOptionsRejectMissingJitter_HP_RUNTIME_02(t *testing.T) {
}
}
func TestRunReportsRetryableSessionError_HP_RUNTIME_04(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
store, err := spool.Open(ctx, spool.Options{DataDir: filepath.Join(t.TempDir(), "spool"), BusyTimeout: time.Second})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = store.Close() })
want := errors.New("dial refused")
reported := make(chan error, 1)
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 = Run(ctx, RunnerOptions{
Store: store, Dial: func(context.Context) (Transport, error) { return nil, want }, Hello: hello,
Limits: agentproto.DefaultLimits(), Backoff: BackoffOptions{Initial: time.Millisecond, Maximum: time.Millisecond, StableReset: time.Second},
Jitter: func(time.Duration) time.Duration { return 0 }, Now: func() time.Time { return time.Now().UTC() },
OnSessionError: func(got error) { reported <- got; cancel() },
})
if err != nil {
t.Fatalf("Run error = %v, want graceful cancellation", err)
}
select {
case got := <-reported:
if !errors.Is(got, want) {
t.Fatalf("reported error = %v, want %v", got, want)
}
default:
t.Fatal("retryable session error was not reported")
}
}
func TestFlushEventsAssignsAndSendsOnlyUnacknowledgedRows_HP_RUNTIME_03(t *testing.T) {
ctx := context.Background()
store, err := spool.Open(ctx, spool.Options{DataDir: filepath.Join(t.TempDir(), "spool"), BusyTimeout: time.Second})
+1 -1
View File
@@ -40,7 +40,7 @@ func defaultServerFile() serverFile {
},
Observability: observabilityFile{
Listen: "127.0.0.1:6901", LivenessPath: "/livez", ReadinessPath: "/readyz",
MetricsPath: "/metrics", LogLevel: "info", LogFormat: "json",
MetricsPath: "/metrics", LogLevel: "info", LogFormat: "json", LogMaxBytes: 10 << 20, LogMaxFiles: 5,
},
}
}
+1 -1
View File
@@ -23,7 +23,7 @@ func TestEffectiveConfigurationGolden_HP_CFG_09(t *testing.T) {
t.Fatal("effective server configuration is nondeterministic")
}
digest := sha256.Sum256(first)
const goldenSHA256 = "55c5fec2c10bc27372353a60f317c5a2959f8087685cf7c11202bb12dc59aa28"
const goldenSHA256 = "65c938efd935ce0c0d2354cdb62c38a2dc177755520d1ae1e5951902a70dd323"
if got := hex.EncodeToString(digest[:]); got != goldenSHA256 {
t.Fatalf("effective server configuration golden changed: got %s", got)
}
+156
View File
@@ -0,0 +1,156 @@
package observability
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"strconv"
"sync"
"time"
)
// RotatingFile is an append-only daemon log with bounded numbered archives.
// It is deliberately a single-process writer: deployment supplies one daemon
// process per configured path, and a second writer must use a distinct file.
// Archive .1 is newest; .maxFiles is oldest.
type RotatingFile struct {
mu sync.Mutex
path string
maxBytes uint64
maxFiles uint32
file *os.File
size uint64
}
// OpenRotatingFile opens path for append, creating its parent directories with
// conservative permissions. A blank path disables file logging and returns a
// no-op closer. Both rotation controls must be zero (unbounded) or positive.
func OpenRotatingFile(path string, maxBytes uint64, maxFiles uint32) (io.WriteCloser, error) {
if path == "" {
return nopWriteCloser{Writer: io.Discard}, nil
}
if (maxBytes == 0) != (maxFiles == 0) {
return nil, errors.New("log rotation byte/file limits must both be zero or positive")
}
if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil {
return nil, fmt.Errorf("create log directory: %w", err)
}
result := &RotatingFile{path: path, maxBytes: maxBytes, maxFiles: maxFiles}
if err := result.open(); err != nil {
return nil, err
}
return result, nil
}
func (file *RotatingFile) open() error {
opened, err := os.OpenFile(file.path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o600)
if err != nil {
return fmt.Errorf("open log file: %w", err)
}
info, err := opened.Stat()
if err != nil {
_ = opened.Close()
return fmt.Errorf("stat log file: %w", err)
}
file.file = opened
file.size = uint64(info.Size())
return nil
}
func (file *RotatingFile) Write(data []byte) (int, error) {
file.mu.Lock()
defer file.mu.Unlock()
if file.file == nil {
return 0, os.ErrClosed
}
if file.maxBytes != 0 && file.size != 0 && (file.size >= file.maxBytes || uint64(len(data)) > file.maxBytes-file.size) {
if err := file.rotate(); err != nil {
return 0, err
}
}
written, err := file.file.Write(data)
file.size += uint64(written)
return written, err
}
func (file *RotatingFile) Close() error {
file.mu.Lock()
defer file.mu.Unlock()
if file.file == nil {
return nil
}
err := file.file.Close()
file.file = nil
return err
}
func (file *RotatingFile) rotate() error {
if err := file.file.Close(); err != nil {
return fmt.Errorf("close log before rotation: %w", err)
}
file.file = nil
oldest := file.archivePath(file.maxFiles)
if err := os.Remove(oldest); err != nil && !errors.Is(err, os.ErrNotExist) {
return fmt.Errorf("remove oldest log archive: %w", err)
}
for index := file.maxFiles; index > 1; index-- {
from := file.archivePath(index - 1)
to := file.archivePath(index)
if err := os.Rename(from, to); err != nil && !errors.Is(err, os.ErrNotExist) {
return fmt.Errorf("rotate log archive: %w", err)
}
}
if err := os.Rename(file.path, file.archivePath(1)); err != nil && !errors.Is(err, os.ErrNotExist) {
return fmt.Errorf("seal current log: %w", err)
}
return file.open()
}
func (file *RotatingFile) archivePath(index uint32) string {
return file.path + "." + strconv.FormatUint(uint64(index), 10)
}
type nopWriteCloser struct{ io.Writer }
func (nopWriteCloser) Close() error { return nil }
// FormatLog makes file records either newline-delimited JSON or plain text.
// It is intentionally applied only to daemon diagnostics, never command
// output, stdin, environment values, or other payload-bearing data.
func FormatLog(destination io.Writer, format string) io.Writer {
if destination == nil || format != "json" {
return destination
}
return &jsonLogWriter{destination: destination}
}
type jsonLogWriter struct {
mu sync.Mutex
destination io.Writer
}
func (writer *jsonLogWriter) Write(data []byte) (int, error) {
writer.mu.Lock()
defer writer.mu.Unlock()
for _, line := range bytes.Split(data, []byte{'\n'}) {
if len(line) == 0 {
continue
}
record, err := json.Marshal(struct {
Time string `json:"time"`
Level string `json:"level"`
Message string `json:"message"`
}{Time: time.Now().UTC().Format(time.RFC3339Nano), Level: "info", Message: string(line)})
if err != nil {
return 0, fmt.Errorf("encode structured log: %w", err)
}
if _, err := writer.destination.Write(append(record, '\n')); err != nil {
return 0, err
}
}
return len(data), nil
}
+86
View File
@@ -0,0 +1,86 @@
package observability
import (
"bytes"
"encoding/json"
"errors"
"os"
"path/filepath"
"strings"
"testing"
)
func TestRotatingFileBoundsArchives_HP_OPS_02(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "nested", "rvbox.log")
writer, err := OpenRotatingFile(path, 5, 2)
if err != nil {
t.Fatal(err)
}
if _, err := writer.Write([]byte("first")); err != nil {
t.Fatal(err)
}
if _, err := writer.Write([]byte("two")); err != nil {
t.Fatal(err)
}
if _, err := writer.Write([]byte("three")); err != nil {
t.Fatal(err)
}
if _, err := writer.Write([]byte("four")); err != nil {
t.Fatal(err)
}
if err := writer.Close(); err != nil {
t.Fatal(err)
}
current, err := os.ReadFile(path)
if err != nil || string(current) != "four" {
t.Fatalf("current = %q, %v", current, err)
}
newest, err := os.ReadFile(path + ".1")
if err != nil || string(newest) != "three" {
t.Fatalf("newest archive = %q, %v", newest, err)
}
oldest, err := os.ReadFile(path + ".2")
if err != nil || string(oldest) != "two" {
t.Fatalf("oldest archive = %q, %v", oldest, err)
}
if _, err := os.Stat(path + ".3"); !errors.Is(err, os.ErrNotExist) {
t.Fatalf("unexpected third archive: %v", err)
}
}
func TestRotatingFileRejectsPartialRotationConfig_BH_OPS_03(t *testing.T) {
t.Parallel()
for _, limits := range [][2]uint64{{1, 0}, {0, 1}} {
_, err := OpenRotatingFile(filepath.Join(t.TempDir(), "rvbox.log"), limits[0], uint32(limits[1]))
if err == nil || !strings.Contains(err.Error(), "both be zero") {
t.Fatalf("limits %v error = %v", limits, err)
}
}
}
func TestFormatLogWritesStructuredBoundedRecords_HP_OPS_04(t *testing.T) {
t.Parallel()
var structured bytes.Buffer
if _, err := FormatLog(&structured, "json").Write([]byte("connected\nrecovered\n")); err != nil {
t.Fatal(err)
}
lines := bytes.Split(bytes.TrimSpace(structured.Bytes()), []byte{'\n'})
if len(lines) != 2 {
t.Fatalf("structured records = %q", structured.String())
}
for index, want := range []string{"connected", "recovered"} {
var record struct {
Time string `json:"time"`
Level string `json:"level"`
Message string `json:"message"`
}
if err := json.Unmarshal(lines[index], &record); err != nil || record.Time == "" || record.Level != "info" || record.Message != want {
t.Fatalf("record %d = %+v, %v", index, record, err)
}
}
var plain bytes.Buffer
if _, err := FormatLog(&plain, "text").Write([]byte("plain\n")); err != nil || plain.String() != "plain\n" {
t.Fatalf("text log = %q, %v", plain.String(), err)
}
}
+26 -14
View File
@@ -55,20 +55,9 @@ func (handler *JSONRPCHandler) ServeHTTP(response http.ResponseWriter, request *
handler.writeRPCError(response, nil, jsonRPCParseError, "could not read request", nil)
return
}
if int64(len(body)) > handler.MaxBody {
handler.writeRPCError(response, nil, jsonRPCInvalid, "request body exceeds limit", nil)
return
}
var envelope jsonRPCRequest
decoder := json.NewDecoder(bytes.NewReader(body))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&envelope); err != nil {
handler.writeRPCError(response, nil, jsonRPCParseError, "invalid JSON", nil)
return
}
var trailing any
if err := decoder.Decode(&trailing); err != io.EOF || envelope.JSONRPC != jsonRPCVersion || envelope.Method == "" || len(envelope.ID) == 0 || bytes.Equal(bytes.TrimSpace(envelope.ID), []byte("null")) {
handler.writeRPCError(response, nil, jsonRPCInvalid, "invalid JSON-RPC request", nil)
envelope, errorCode, errorMessage := decodeJSONRPCRequest(body, handler.MaxBody)
if errorCode != 0 {
handler.writeRPCError(response, nil, errorCode, errorMessage, nil)
return
}
result, callErr := handler.call(request.Context(), envelope.Method, envelope.Params)
@@ -87,6 +76,29 @@ type jsonRPCRequest struct {
Params json.RawMessage `json:"params"`
}
// decodeJSONRPCRequest keeps the hostile JSON boundary independently bounded
// and fuzzable. The returned code/message are the externally stable JSON-RPC
// parse or invalid-request result; callers must not inspect partial fields.
func decodeJSONRPCRequest(body []byte, maxBody int64) (jsonRPCRequest, int, string) {
if maxBody <= 0 {
maxBody = defaultJSONRPCBody
}
if int64(len(body)) > maxBody {
return jsonRPCRequest{}, jsonRPCInvalid, "request body exceeds limit"
}
var envelope jsonRPCRequest
decoder := json.NewDecoder(bytes.NewReader(body))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&envelope); err != nil {
return jsonRPCRequest{}, jsonRPCParseError, "invalid JSON"
}
var trailing any
if err := decoder.Decode(&trailing); err != io.EOF || envelope.JSONRPC != jsonRPCVersion || envelope.Method == "" || len(envelope.ID) == 0 || bytes.Equal(bytes.TrimSpace(envelope.ID), []byte("null")) {
return jsonRPCRequest{}, jsonRPCInvalid, "invalid JSON-RPC request"
}
return envelope, 0, ""
}
type jsonRPCResponse struct {
JSONRPC string `json:"jsonrpc"`
ID json.RawMessage `json:"id"`
+15
View File
@@ -95,3 +95,18 @@ func TestJSONRPCRejectsOversizeAndMalformedRequests_BH_CTL_16(t *testing.T) {
t.Fatalf("malformed response = %s", malformedBody)
}
}
func FuzzDecodeJSONRPCRequestBounded_SEC_CTL_01(f *testing.F) {
f.Add([]byte(`{"jsonrpc":"2.0","id":1,"method":"listClients","params":{}}`))
f.Add([]byte(`{"jsonrpc":"2.0","id":null,"method":"listClients"}`))
f.Add([]byte(`{"jsonrpc":"2.0","id":1,"method":"listClients"}{}`))
f.Fuzz(func(t *testing.T, body []byte) {
// Keep fuzzing at the same independently enforced boundary as the
// production handler rather than allowing a corpus entry to allocate
// unbounded JSON decoder state.
if len(body) > 64<<10 {
return
}
_, _, _ = decodeJSONRPCRequest(body, 64<<10)
})
}
+46 -10
View File
@@ -102,6 +102,10 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
return
}
hello := helloEnvelope.GetClientHello()
// This cutoff precedes registration, which is when control callers can
// first observe the live client. Reconciliation must compare only the
// state included in its request; later control commands are fresh work.
reconcileBoundary := server.now()
instanceID, err := domain.ParseUUIDv7(hello.GetClientInstanceId())
if err != nil {
server.close(connection, websocket.StatusPolicyViolation, "invalid client instance ID")
@@ -176,7 +180,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
server.close(connection, websocket.StatusPolicyViolation, "invalid initial client capacity")
return
}
go server.dispatchLoop(sessionContext, queue, handle.DispatchWake(), reconciled, &capacityMu, capacity, reservations, stdinSent, signalSent, scriptTransfers, hello.GetClientId(), hello.GetPlatform(), encodeSessionID(sessionID), registration.Generation, cancel)
go server.dispatchLoop(sessionContext, queue, handle.DispatchWake(), handle.SignalDispatch, reconciled, &capacityMu, capacity, reservations, stdinSent, signalSent, scriptTransfers, hello.GetClientId(), hello.GetPlatform(), encodeSessionID(sessionID), registration.Generation, cancel)
encodedSessionID := encodeSessionID(sessionID)
welcome, err := proto.Marshal(&rvboxv1.AgentEnvelope{
SessionId: encodedSessionID, SessionGeneration: registration.Generation,
@@ -192,7 +196,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
server.close(connection, websocket.StatusInternalError, "could not queue welcome")
return
}
targets, err := server.Store.ReconcileTargets(parent, hello.GetClientId())
targets, err := server.Store.ReconcileTargetsAt(parent, hello.GetClientId(), reconcileBoundary)
if err != nil {
server.close(connection, websocket.StatusInternalError, "could not build reconciliation request")
return
@@ -237,7 +241,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
return
}
if snapshot := envelope.GetReconcileSnapshot(); snapshot != nil {
result, reconcileErr := server.Store.ReconcileClientSnapshotForSession(sessionContext, hello.GetClientId(), registration.Generation, snapshot)
result, reconcileErr := server.Store.ReconcileClientSnapshotForSessionAt(sessionContext, hello.GetClientId(), registration.Generation, snapshot, reconcileBoundary)
if reconcileErr != nil {
server.close(connection, websocket.StatusPolicyViolation, "reconciliation failed")
return
@@ -409,7 +413,7 @@ func eventType(event *rvboxv1.CommandEvent) uint16 {
// enqueueNextDispatch records the queued-to-dispatched transition before
// exposing work to the network. A full data lane is a pre-write failure, so
// only the owning generation can put the command back into the queue.
func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *WriterQueue, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, beforeEnqueue func(*store.DispatchCandidate)) (domain.UUID, bool, error) {
func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *WriterQueue, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, beforeEnqueue func(*store.DispatchCandidate), onWritten func(domain.UUID)) (domain.UUID, bool, error) {
candidate, err := server.Store.ClaimNextDispatch(ctx, clientID, generation, server.now())
if err != nil || candidate == nil {
return domain.UUID{}, false, err
@@ -447,7 +451,11 @@ func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *Write
if beforeEnqueue != nil {
beforeEnqueue(candidate)
}
if !queue.EnqueueData(Frame{Kind: FrameData, Payload: encoded}) {
if !queue.EnqueueData(Frame{Kind: FrameData, Payload: encoded, OnWritten: func() {
if onWritten != nil {
onWritten(candidate.IssueUUID)
}
}}) {
sent, requeueErr := requeue(ErrDispatchDataFull)
return candidate.IssueUUID, sent, requeueErr
}
@@ -527,12 +535,18 @@ func (server *AgentServer) enqueueNextScript(queue *WriterQueue, sessionID strin
// lane. Session-local sent tracking suppresses duplicate frames while a live
// connection remains usable; reconnecting naturally replays unacknowledged
// writes from storage.
func (server *AgentServer) enqueueNextStdin(ctx context.Context, queue *WriterQueue, clientID, sessionID string, generation uint64, sent map[string]struct{}, sentMu *sync.Mutex) (bool, error) {
func (server *AgentServer) enqueueNextStdin(ctx context.Context, queue *WriterQueue, clientID, sessionID string, generation uint64, sent map[string]struct{}, dispatchWritten map[string]bool, sentMu *sync.Mutex) (bool, error) {
intents, err := server.Store.PendingStdin(ctx, clientID)
if err != nil {
return false, err
}
for _, intent := range intents {
sentMu.Lock()
written, waitingForDispatch := dispatchWritten[intent.IssueUUID.String()]
sentMu.Unlock()
if waitingForDispatch && !written {
continue
}
key := stdinIntentKey(intent.IssueUUID, intent.WriteSeq)
sentMu.Lock()
_, alreadySent := sent[key]
@@ -575,12 +589,18 @@ func signalIntentKey(issue domain.UUID, revision uint64, signal rvboxv1.SignalKi
return issue.String() + ":" + fmt.Sprint(revision) + ":" + fmt.Sprint(signal)
}
func (server *AgentServer) enqueueNextSignal(ctx context.Context, queue *WriterQueue, clientID, sessionID string, generation uint64, sent map[string]struct{}, sentMu *sync.Mutex) (bool, error) {
func (server *AgentServer) enqueueNextSignal(ctx context.Context, queue *WriterQueue, clientID, sessionID string, generation uint64, sent map[string]struct{}, dispatchWritten map[string]bool, sentMu *sync.Mutex) (bool, error) {
intents, err := server.Store.PendingSignals(ctx, clientID, generation)
if err != nil {
return false, err
}
for _, intent := range intents {
sentMu.Lock()
written, waitingForDispatch := dispatchWritten[intent.IssueUUID.String()]
sentMu.Unlock()
if waitingForDispatch && !written {
continue
}
key := signalIntentKey(intent.IssueUUID, intent.CommandRevision, intent.Signal)
sentMu.Lock()
_, alreadySent := sent[key]
@@ -613,7 +633,11 @@ func (server *AgentServer) enqueueNextSignal(ctx context.Context, queue *WriterQ
// dispatchLoop is the per-session serialized dispatcher. It waits for a
// complete reconciliation result before consuming queued work, then coalesces
// wakeups from local control RPCs, capacity advertisements, and acceptances.
func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue, wake <-chan struct{}, reconciled <-chan struct{}, capacityMu *sync.Mutex, capacity *CapacityShadow, reservations map[string]DispatchLane, stdinSent, signalSent map[string]struct{}, scriptTransfers map[string]*scriptTransfer, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, cancel context.CancelFunc) {
func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue, wake <-chan struct{}, signalWake func(), reconciled <-chan struct{}, capacityMu *sync.Mutex, capacity *CapacityShadow, reservations map[string]DispatchLane, stdinSent, signalSent map[string]struct{}, scriptTransfers map[string]*scriptTransfer, clientID string, platform rvboxv1.Platform, sessionID string, generation uint64, cancel context.CancelFunc) {
// A control frame may use the essential queue and therefore overtake data.
// Keep its issue blocked only until the dispatch frame has actually crossed
// the socket writer; after that, WebSocket ordering preserves the dependency.
dispatchWritten := make(map[string]bool)
select {
case <-reconciled:
case <-ctx.Done():
@@ -621,7 +645,7 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue,
}
for {
for {
stdinQueued, stdinErr := server.enqueueNextStdin(ctx, queue, clientID, sessionID, generation, stdinSent, capacityMu)
stdinQueued, stdinErr := server.enqueueNextStdin(ctx, queue, clientID, sessionID, generation, stdinSent, dispatchWritten, capacityMu)
if stdinErr != nil {
server.closeForDispatchFailure(cancel)
return
@@ -629,7 +653,7 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue,
if stdinQueued {
continue
}
signalQueued, signalErr := server.enqueueNextSignal(ctx, queue, clientID, sessionID, generation, signalSent, capacityMu)
signalQueued, signalErr := server.enqueueNextSignal(ctx, queue, clientID, sessionID, generation, signalSent, dispatchWritten, capacityMu)
if signalErr != nil {
server.closeForDispatchFailure(cancel)
return
@@ -655,16 +679,25 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue,
issue, sent, err := server.enqueueNextDispatch(ctx, queue, clientID, platform, sessionID, generation, func(candidate *store.DispatchCandidate) {
capacityMu.Lock()
reservations[candidate.IssueUUID.String()] = lane
dispatchWritten[candidate.IssueUUID.String()] = false
if candidate.ScriptPresent {
scriptTransfers[candidate.IssueUUID.String()] = &scriptTransfer{Body: append([]byte(nil), candidate.ScriptContent...), Digest: sha256.Sum256(candidate.ScriptContent)}
}
inserted = true
capacityMu.Unlock()
}, func(issue domain.UUID) {
capacityMu.Lock()
dispatchWritten[issue.String()] = true
capacityMu.Unlock()
if signalWake != nil {
signalWake()
}
})
if !sent {
capacityMu.Lock()
if inserted {
delete(reservations, issue.String())
delete(dispatchWritten, issue.String())
delete(scriptTransfers, issue.String())
}
capacity.Release(lane)
@@ -706,6 +739,9 @@ func (server *AgentServer) writeLoop(ctx context.Context, connection *websocket.
if frame.Written != nil {
close(frame.Written)
}
if frame.OnWritten != nil {
frame.OnWritten()
}
continue
}
if !errors.Is(err, context.DeadlineExceeded) {
+57 -1
View File
@@ -98,6 +98,57 @@ func TestAgentServerRegistrationAndReplacement_HP_SES_05(t *testing.T) {
}
}
func TestAgentServerDispatchesCommandAdmittedDuringReconcile_HP_SES_11(t *testing.T) {
server, agent, cleanup := newTestAgentServerWithStore(t)
defer cleanup()
instance, err := domain.NewUUIDv7()
if err != nil {
t.Fatal(err)
}
connection, welcome := dialAndHello(t, server, "client-race", instance.String())
defer connection.CloseNow()
issue, err := domain.NewUUIDv7()
if err != nil {
t.Fatal(err)
}
spec, err := proto.Marshal(&rvboxv1.ExecutionSpec{ShellType: rvboxv1.ShellType_SHELL_CMD, Source: &rvboxv1.ExecutionSpec_CommandText{CommandText: "echo dispatch"}})
if err != nil {
t.Fatal(err)
}
now := time.Now().UTC()
if _, err := agent.Store.QueueCommand(context.Background(), store.QueueCommandInput{IssueUUID: issue, ClientID: "client-race", IssueTime: now, ReceiptTime: now, ImmutableSHA256: sha256.Sum256([]byte("during-reconcile")), ExecutionSpec: spec}); err != nil {
t.Fatal(err)
}
snapshot, err := proto.Marshal(&rvboxv1.AgentEnvelope{
SessionId: welcome.GetSessionId(), SessionGeneration: welcome.GetSessionGeneration(),
Payload: &rvboxv1.AgentEnvelope_ReconcileSnapshot{ReconcileSnapshot: &rvboxv1.ReconcileSnapshot{}},
})
if err != nil {
t.Fatal(err)
}
if err := connection.Write(context.Background(), websocket.MessageBinary, snapshot); err != nil {
t.Fatal(err)
}
readContext, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
for index := 0; index < 2; index++ {
_, payload, err := connection.Read(readContext)
if err != nil {
t.Fatal(err)
}
var envelope rvboxv1.AgentEnvelope
if err := proto.Unmarshal(payload, &envelope); err != nil {
t.Fatal(err)
}
if index == 0 && envelope.GetReconcileResult() == nil {
t.Fatalf("first post-snapshot envelope = %T, want reconcile result", envelope.Payload)
}
if index == 1 && (envelope.GetCommandDispatch() == nil || envelope.GetCommandDispatch().GetIssueUuid() != issue.String()) {
t.Fatalf("second post-snapshot envelope = %+v, want dispatch %s", envelope.Payload, issue)
}
}
}
func TestAgentServerRejectsHostileWireInputs_BH_SES_04(t *testing.T) {
server, cleanup := newTestAgentServer(t)
defer cleanup()
@@ -174,6 +225,11 @@ func TestWireEventAppendCarriesClientBinding_HP_EVENT_01(t *testing.T) {
}
func newTestAgentServer(t *testing.T) (*httptest.Server, func()) {
server, _, cleanup := newTestAgentServerWithStore(t)
return server, cleanup
}
func newTestAgentServerWithStore(t *testing.T) (*httptest.Server, *AgentServer, func()) {
t.Helper()
dataDirectory := filepath.Join(t.TempDir(), "store")
if err := os.Mkdir(dataDirectory, 0o700); err != nil {
@@ -190,7 +246,7 @@ func newTestAgentServer(t *testing.T) (*httptest.Server, func()) {
HeartbeatIdle: time.Second, LivenessTimeout: 2 * time.Second,
}
httpServer := httptest.NewServer(agent)
return httpServer, func() {
return httpServer, agent, func() {
httpServer.Close()
if err := persistence.Close(); err != nil {
t.Errorf("close persistence: %v", err)
+4
View File
@@ -26,6 +26,10 @@ type Frame struct {
// Written is closed by the sole socket writer after a successful write.
// It is used only for protocol barriers such as reconciliation-before-work.
Written chan<- struct{}
// OnWritten is a non-blocking session-local scheduler notification. It is
// invoked only after the sole socket writer has completed the frame, so a
// dependent control frame cannot overtake its command dispatch.
OnWritten func()
}
// WriterQueue is owned by one socket writer. Data saturation leaves the work
+31
View File
@@ -126,6 +126,37 @@ func TestClaimDispatchExpiresAndFencesRequeue_HP_DISPATCH_02(t *testing.T) {
}
}
func TestPendingStdinWaitsForDurableDispatch_HP_DISPATCH_06(t *testing.T) {
t.Parallel()
ctx := context.Background()
opened, err := Open(ctx, Options{DataDir: filepath.Join(t.TempDir(), "state"), BusyTimeout: time.Second})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = opened.Close() })
now := time.Date(2026, time.September, 11, 8, 50, 0, 0, time.UTC)
if _, err := opened.RegisterClientSession(ctx, ClientRegistration{ClientID: "stdin-client", Platform: 2, Architecture: "amd64", DaemonVersion: "test", DaemonCWD: `C:\\`, SupportedShells: []byte{1}, ClientInstanceID: [16]byte{41}, SessionID: [16]byte{42}, ConnectedAt: now}); err != nil {
t.Fatal(err)
}
issue := fixedStoreIssue(0x97)
if _, err := opened.QueueCommand(ctx, QueueCommandInput{IssueUUID: issue, ClientID: "stdin-client", IssueTime: now, ReceiptTime: now, ImmutableSHA256: sha256.Sum256([]byte("stdin ordering")), ExecutionSpec: []byte("spec")}); err != nil {
t.Fatal(err)
}
if _, err := opened.AppendStdin(ctx, StdinWriteInput{IssueUUID: issue, ClientID: "stdin-client", RequestUUID: fixedStoreIssue(0x98), Data: []byte("input"), ImmutableHash: sha256.Sum256([]byte("stdin append")), OccurredAt: now}); err != nil {
t.Fatal(err)
}
if pending, err := opened.PendingStdin(ctx, "stdin-client"); err != nil || len(pending) != 0 {
t.Fatalf("queued command exposed stdin = %#v, %v", pending, err)
}
if candidate, err := opened.ClaimNextDispatch(ctx, "stdin-client", 7, now); err != nil || candidate == nil || candidate.IssueUUID != issue {
t.Fatalf("dispatch claim = %#v, %v", candidate, err)
}
pending, err := opened.PendingStdin(ctx, "stdin-client")
if err != nil || len(pending) != 1 || pending[0].IssueUUID != issue || pending[0].WriteSeq != 1 || string(pending[0].Data) != "input" {
t.Fatalf("dispatched command stdin = %#v, %v", pending, err)
}
}
func TestLateDispatchAcceptanceAndEventRetainExpiryContradiction_BH_DISPATCH_09(t *testing.T) {
t.Parallel()
ctx := context.Background()
+40 -6
View File
@@ -16,11 +16,26 @@ import (
// client. It is a read-only snapshot used to tell a reconnecting agent which
// UUIDs and durable cursors must be compared before fresh dispatch is enabled.
func (store *Store) ReconcileTargets(ctx context.Context, clientID string) ([]*rvboxv1.ReconcileTarget, error) {
return store.ReconcileTargetsAt(ctx, clientID, time.Time{})
}
// ReconcileTargetsAt returns the durable non-terminal view at the supplied
// receipt-time boundary. A connecting session uses one boundary for both the
// request and its reply: commands admitted after it are fresh work, not absent
// client evidence to be reconciled.
func (store *Store) ReconcileTargetsAt(ctx context.Context, clientID string, receiptBoundary time.Time) ([]*rvboxv1.ReconcileTarget, error) {
if clientID == "" {
return nil, errors.New("client ID is required")
}
rows, err := store.db.QueryContext(ctx, `SELECT issue_uuid, last_event_seq, revision, immutable_request_sha256
FROM commands WHERE client_id = ? AND lifecycle BETWEEN 1 AND 4 ORDER BY issue_time, issue_uuid`, clientID)
query := `SELECT issue_uuid, last_event_seq, revision, immutable_request_sha256
FROM commands WHERE client_id = ? AND lifecycle BETWEEN 1 AND 4`
args := []any{clientID}
if !receiptBoundary.IsZero() {
query += ` AND server_receipt_time <= ?`
args = append(args, receiptBoundary.UTC().UnixNano())
}
query += ` ORDER BY issue_time, issue_uuid`
rows, err := store.db.QueryContext(ctx, query, args...)
if err != nil {
return nil, err
}
@@ -62,10 +77,19 @@ func (store *Store) ReconcileClientSnapshot(ctx context.Context, clientID string
// incidented. Retained non-terminal rows are retargeted to this generation so
// late events from the previous connection cannot advance the command.
func (store *Store) ReconcileClientSnapshotForSession(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot) (*rvboxv1.ReconcileResult, error) {
return store.ReconcileClientSnapshotForSessionAt(ctx, clientID, generation, snapshot, time.Time{})
}
// ReconcileClientSnapshotForSessionAt reconciles exactly the server state that
// was included in the matching ReconcileRequest. Work admitted after the
// boundary is intentionally left for the dispatch loop once reconciliation
// completes; treating it as absent client evidence would lose a valid command
// during the Hello/reconcile race.
func (store *Store) ReconcileClientSnapshotForSessionAt(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot, receiptBoundary time.Time) (*rvboxv1.ReconcileResult, error) {
if generation == 0 {
return nil, errors.New("session generation is required")
}
return store.reconcileClientSnapshot(ctx, clientID, generation, snapshot)
return store.reconcileClientSnapshotAt(ctx, clientID, generation, snapshot, receiptBoundary)
}
type reconcileServerRow struct {
@@ -78,6 +102,10 @@ type reconcileServerRow struct {
}
func (store *Store) reconcileClientSnapshot(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot) (*rvboxv1.ReconcileResult, error) {
return store.reconcileClientSnapshotAt(ctx, clientID, generation, snapshot, time.Time{})
}
func (store *Store) reconcileClientSnapshotAt(ctx context.Context, clientID string, generation uint64, snapshot *rvboxv1.ReconcileSnapshot, receiptBoundary time.Time) (*rvboxv1.ReconcileResult, error) {
if clientID == "" {
return nil, errors.New("client ID is required")
}
@@ -112,7 +140,7 @@ func (store *Store) reconcileClientSnapshot(ctx context.Context, clientID string
return nil, err
}
defer tx.Rollback()
serverRows, tombstones, err := loadReconcileRows(ctx, tx, clientID)
serverRows, tombstones, err := loadReconcileRows(ctx, tx, clientID, receiptBoundary)
if err != nil {
store.writeMu.Unlock()
return nil, err
@@ -215,8 +243,14 @@ type reconcileIncident struct {
dataLoss bool
}
func loadReconcileRows(ctx context.Context, tx *sql.Tx, clientID string) (map[string]reconcileServerRow, map[string][]byte, error) {
rows, err := tx.QueryContext(ctx, `SELECT issue_uuid, lifecycle, revision, last_event_seq, immutable_request_sha256, target_session_generation FROM commands WHERE client_id = ?`, clientID)
func loadReconcileRows(ctx context.Context, tx *sql.Tx, clientID string, receiptBoundary time.Time) (map[string]reconcileServerRow, map[string][]byte, error) {
query := `SELECT issue_uuid, lifecycle, revision, last_event_seq, immutable_request_sha256, target_session_generation FROM commands WHERE client_id = ?`
args := []any{clientID}
if !receiptBoundary.IsZero() {
query += ` AND server_receipt_time <= ?`
args = append(args, receiptBoundary.UTC().UnixNano())
}
rows, err := tx.QueryContext(ctx, query, args...)
if err != nil {
return nil, nil, err
}
+25
View File
@@ -71,6 +71,31 @@ func TestReconcileSnapshotMutatesMissingAndRetargets_HP_SES_13(t *testing.T) {
}
}
func TestReconcileSessionBoundaryLeavesNewlyQueuedCommandForDispatch_HP_RECONCILE_08(t *testing.T) {
t.Parallel()
ctx := context.Background()
opened, err := Open(ctx, Options{DataDir: filepath.Join(t.TempDir(), "state"), BusyTimeout: time.Second})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = opened.Close() })
boundary := time.Date(2026, time.September, 11, 12, 0, 0, 0, time.UTC)
if _, err := opened.RegisterClientSession(ctx, ClientRegistration{ClientID: "boundary-client", Platform: 3, Architecture: "amd64", DaemonVersion: "test", DaemonCWD: `C:\\`, SupportedShells: []byte{1}, ClientInstanceID: [16]byte{31}, SessionID: [16]byte{32}, ConnectedAt: boundary}); err != nil {
t.Fatal(err)
}
issue := mustReconcileIssue(t, "019c46f1-1d02-7000-8000-0000000000d1")
if _, err := opened.QueueCommand(ctx, QueueCommandInput{IssueUUID: issue, ClientID: "boundary-client", IssueTime: boundary.Add(time.Nanosecond), ReceiptTime: boundary.Add(time.Nanosecond), ImmutableSHA256: sha256.Sum256([]byte("new-after-reconcile-boundary")), ExecutionSpec: []byte("echo queued")}); err != nil {
t.Fatal(err)
}
if _, err := opened.ReconcileClientSnapshotForSessionAt(ctx, "boundary-client", 1, &rvboxv1.ReconcileSnapshot{}, boundary); err != nil {
t.Fatal(err)
}
view, err := opened.GetCommandView(ctx, "boundary-client", issue)
if err != nil || view.Lifecycle != uint32(rvboxv1.CommandLifecycle_COMMAND_QUEUED) {
t.Fatalf("post-boundary command = %#v, %v", view, err)
}
}
func mustReconcileIssue(t *testing.T, value string) domain.UUID {
t.Helper()
issue, err := domain.ParseUUIDv7(value)
+7 -5
View File
@@ -46,10 +46,12 @@ type StdinIntent struct {
Close bool
}
// PendingStdin returns unacknowledged input intents for one client. Delivery is
// intentionally tracked by the session dispatcher, not by this durable query:
// a lost connection simply causes the next session to replay the same write
// sequence, which the client acknowledges idempotently.
// PendingStdin returns unacknowledged input intents only after the command has
// entered DISPATCHED state. This durable gate prevents append/close controls
// from overtaking their initial CommandDispatch on a fresh session. Delivery is
// otherwise tracked by the session dispatcher: a lost connection simply causes
// the next session to replay the same write sequence, which the client
// acknowledges idempotently.
func (store *Store) PendingStdin(ctx context.Context, clientID string) ([]StdinIntent, error) {
if clientID == "" {
return nil, errors.New("client ID is required")
@@ -60,7 +62,7 @@ func (store *Store) PendingStdin(ctx context.Context, clientID string) ([]StdinI
}
rows, err := database.QueryContext(ctx, `SELECT w.issue_uuid, w.write_seq, w.payload, w.raw_bytes, w.stored_bytes, w.compression, w.sha256, w.append_newline, w.close_intent
FROM stdin_writes w JOIN commands c ON c.issue_uuid = w.issue_uuid
WHERE c.client_id = ? AND c.lifecycle BETWEEN 1 AND 4 AND w.acknowledged = 0
WHERE c.client_id = ? AND c.lifecycle BETWEEN 2 AND 4 AND w.acknowledged = 0
ORDER BY c.issue_time, c.issue_uuid, w.write_seq`, clientID)
if err != nil {
return nil, err
+39
View File
@@ -0,0 +1,39 @@
#!/bin/sh
# Run one bounded Go fuzz target in the pinned RVBox toolchain container.
set -eu
repo_root=$(CDPATH= cd -- "$(dirname -- "$0")/.." && pwd)
usage() {
cat <<'EOF'
usage: scripts/test-fuzz --package ./PACKAGE --name FUZZ_TEST [--time DURATION]
Runs exactly one Fuzz* target with ordinary unit tests disabled. The default
duration is 30s. A crash leaves Go's minimal reproducer in the package fuzz
corpus, where it becomes a normal deterministic seed on the next test run.
EOF
}
fail() { printf '%s\n' "test-fuzz: $*" >&2; exit 2; }
package=
name=
duration=30s
while [ "$#" -gt 0 ]; do
case $1 in
--package) [ "$#" -ge 2 ] || fail "--package needs a value"; package=$2; shift 2 ;;
--name) [ "$#" -ge 2 ] || fail "--name needs a value"; name=$2; shift 2 ;;
--time) [ "$#" -ge 2 ] || fail "--time needs a value"; duration=$2; shift 2 ;;
--help|-h) usage; exit 0 ;;
*) fail "unknown argument $1" ;;
esac
done
case $package in ./*) ;; *) fail "--package must be a repository-relative Go package" ;; esac
case $name in Fuzz*) ;; *) fail "--name must be a Fuzz* test function" ;; esac
case $name in *[!A-Za-z0-9_]* ) fail "--name contains unsupported characters" ;; esac
case $duration in *[!0-9a-zA-Z.]*) fail "--time contains unsupported characters" ;; esac
cd "$repo_root"
exec docker compose -f deploy/compose.yaml run --rm toolchain \
go test "$package" -run '^$' -fuzz "$name" -fuzztime "$duration"
+75
View File
@@ -78,6 +78,15 @@ compose() {
docker compose -p "$project" -f "$repo_root/test/linux-server/compose.yaml" "$@"
}
remove_project_networks() {
# The project label is Docker Compose's exact resource-ownership key. This
# also reclaims legacy certgen default networks left by fixtures created
# before certgen joined the labeled native network.
for network_id in $(docker network ls --filter "label=com.docker.compose.project=$project" --quiet); do
docker network rm "$network_id" >/dev/null || fail "could not remove owned Docker network $network_id"
done
}
rvc() {
compose exec -T server /opt/rvbox/rvc --socket /run/rvbox/server.sock "$@"
}
@@ -179,6 +188,69 @@ assert_stdin_close() {
fail "stdin command did not reach terminal success"
}
assert_signal_term() {
issued=$(rvc run --background --shell cmd "$client_id" 'ping -t 127.0.0.1 >NUL' 2>&1) || fail "signal command admission failed: $issued"
issue=$(printf '%s\n' "$issued" | awk 'NR == 1 { print $1 }')
case $issue in ????????-????-7???-????-????????????) ;; *) fail "signal command returned invalid issue UUID: $issued" ;; esac
attempt=0
while [ "$attempt" -lt 30 ]; do
result=$(rvc stat "$client_id" "$issue" 2>/dev/null || true)
if printf '%s\n' "$result" | grep -q 'lifecycle=COMMAND_RUNNING'; then break; fi
attempt=$((attempt + 1)); sleep 1
done
[ "$attempt" -lt 30 ] || fail "signal command did not reach running state"
rvc kill TERM "$client_id" "$issue" >/dev/null || fail "TERM request failed"
attempt=0
while [ "$attempt" -lt 45 ]; do
result=$(rvc stat "$client_id" "$issue" 2>/dev/null || true)
if printf '%s\n' "$result" | grep -q 'lifecycle=COMMAND_TERMINATED'; then
printf '%s\n' "$result" >"$fixture_dir/signal-term.stat"
printf 'passed signal-term issue=%s\n' "$issue"
return 0
fi
case $result in *'lifecycle=COMMAND_FAILED'*|*'lifecycle=COMMAND_REJECTED'*|*'lifecycle=COMMAND_SUCCEEDED'*) fail "TERM command reached wrong terminal state: $result" ;; esac
attempt=$((attempt + 1)); sleep 1
done
fail "TERM command did not reach terminal state"
}
assert_server_restart_reconnect() {
issued=$(rvc run --background --shell cmd "$client_id" 'ping -t 127.0.0.1 >NUL' 2>&1) || fail "restart command admission failed: $issued"
issue=$(printf '%s\n' "$issued" | awk 'NR == 1 { print $1 }')
case $issue in ????????-????-7???-????-????????????) ;; *) fail "restart command returned invalid issue UUID: $issued" ;; esac
attempt=0
while [ "$attempt" -lt 30 ]; do
result=$(rvc stat "$client_id" "$issue" 2>/dev/null || true)
if printf '%s\n' "$result" | grep -q 'lifecycle=COMMAND_RUNNING'; then break; fi
attempt=$((attempt + 1)); sleep 1
done
[ "$attempt" -lt 30 ] || fail "restart command did not reach running state"
# This preserves the server state volume while replacing the actual server
# process behind nginx. A later TERM terminal result can only arrive after
# the Windows daemon has established a fresh reconciled WSS session.
compose restart server
attempt=0
while [ "$attempt" -lt 45 ]; do
if rvc stat "$client_id" >/dev/null 2>&1; then break; fi
attempt=$((attempt + 1)); sleep 1
done
[ "$attempt" -lt 45 ] || fail "server control socket did not recover after restart"
wait_client
rvc kill TERM "$client_id" "$issue" >/dev/null || fail "TERM after server restart failed"
attempt=0
while [ "$attempt" -lt 45 ]; do
result=$(rvc stat "$client_id" "$issue" 2>/dev/null || true)
if printf '%s\n' "$result" | grep -q 'lifecycle=COMMAND_TERMINATED'; then
printf '%s\n' "$result" >"$fixture_dir/server-restart-reconnect.stat"
printf 'passed server-restart-reconnect issue=%s\n' "$issue"
return 0
fi
case $result in *'lifecycle=COMMAND_FAILED'*|*'lifecycle=COMMAND_REJECTED'*|*'lifecycle=COMMAND_SUCCEEDED'*) fail "restart TERM command reached wrong terminal state: $result" ;; esac
attempt=$((attempt + 1)); sleep 1
done
fail "restart TERM command did not reach terminal state"
}
collect() {
if [ "$vm_prepared" = yes ]; then
"$repo_root/scripts/windows/test-host" collect --run-id "$run_id" || true
@@ -190,6 +262,7 @@ collect() {
clean() {
compose down --volumes --remove-orphans || true
remove_project_networks
# A reset is the isolation boundary for the next run. Do not conceal a
# failed shutdown/snapshot restore behind a successful-looking `clean`:
# callers must repair or explicitly inspect the retained VM lease first.
@@ -224,6 +297,8 @@ case $action in
"$repo_root/scripts/windows/test-host" run --run-id "$run_id" --endpoint "$endpoint_host:$port"
wait_client
assert_stdin_close
assert_signal_term
assert_server_restart_reconnect
assert_context active-user no active-user
assert_context active-user-elevated yes active-user-elevated
"$repo_root/scripts/windows/test-host" run --run-id "$run_id" --fail-contexts ACTIVE_USER_ELEVATED
+34 -9
View File
@@ -23,6 +23,7 @@ Actions:
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
stop stop RVBox through SCM and request a graceful guest shutdown
reset stop the guest if necessary, restore the declared baseline, and leave it off
recover read-only fixture/run-state check for a stopped-resumable run
@@ -87,7 +88,7 @@ while [ "$#" -gt 0 ]; do
esac
done
case $action in status|prepare|stage|install|run|logoff|collect|inspect|stop|reset|recover) ;; *) usage >&2; fail "unknown action $action" ;; esac
case $action in status|prepare|stage|install|run|logoff|collect|inspect|logs|stop|reset|recover) ;; *) usage >&2; fail "unknown action $action" ;; esac
if [ "$action" != status ]; then
[ -n "$run_id" ] || fail "$action requires --run-id"
safe_id "$run_id"
@@ -169,7 +170,7 @@ remote() {
# retry. Lifecycle transitions remain single-attempt: their caller must
# inspect/recover rather than risk a duplicate reset, shutdown, or logoff.
case $remote_action in
status|recover|prepare-stage|stage|stage-create-root|stage-copy-exe|stage-copy-config|stage-copy-ca|collect|inspect|install|run)
status|recover|prepare-stage|stage|stage-create-root|stage-copy-exe|stage-copy-config|stage-copy-ca|collect|inspect|logs|install|run)
retry_limit=4
;;
esac
@@ -488,9 +489,9 @@ case "$action" in
require_lease
install -d -m 700 "$host_stage/artifacts"
if [ "$(state)" = running ]; then
guest_run --exe 'C:\Windows\System32\cmd.exe' --wait-stdout --wait-stderr --unquoted-args -- \
/d /s /c "sc.exe queryex RVBoxClient > \"$guest_root\\service-status.txt\" 2>&1 & echo RVBOX_GUEST_OK" >/dev/null || true
VBoxManage guestcontrol "$vm" --username "$guest_user" --passwordfile "$password_file" \
provisioner_run --exe 'C:\Windows\System32\cmd.exe' --wait-stdout --wait-stderr --unquoted-args -- \
/d /s /c "sc.exe queryex RVBoxClient > \"$guest_root\\service-status.txt\" 2>&1 & if exist \"C:/ProgramData/RVBox/service-startup.log\" copy /y \"C:/ProgramData/RVBox/service-startup.log\" \"$guest_root\\service-startup.log\" >NUL & if exist \"C:/ProgramData/RVBox/test-logs/rvbox.log\" copy /y \"C:/ProgramData/RVBox/test-logs/rvbox.log\" \"$guest_root\\rvbox.log\" >NUL & echo RVBOX_GUEST_OK" >/dev/null || true
VBoxManage guestcontrol "$vm" --username "$provisioner_user" --passwordfile "$provisioner_password_file" \
copyfrom "$guest_root" "$host_stage/artifacts" --recursive </dev/null >/dev/null 2>&1 || true
fi
printf 'collected host_stage=%s/artifacts\n' "$host_stage"
@@ -500,7 +501,14 @@ case "$action" in
require_lease
[ "$(state)" = running ] || fail "inspect requires a running prepared VM"
guest_run --exe 'C:\Windows\System32\cmd.exe' --wait-stdout --wait-stderr --unquoted-args -- \
/d /s /c "sc.exe queryex RVBoxClient & sc.exe qc RVBoxClient & reg.exe query \"HKLM\\SYSTEM\\CurrentControlSet\\Services\\RVBoxClient\" /v ImagePath & reg.exe query \"HKLM\\SYSTEM\\CurrentControlSet\\Services\\RVBoxClient\" /v ObjectName & dir \"C:/ProgramData/RVBox\" & icacls \"C:/ProgramData/RVBox\" & certutil -hashfile \"$guest_root\\rvbox.exe\" SHA256 & \"$guest_root\\rvbox.exe\" --check-config --config \"$guest_root\\client.toml\" > \"$guest_root\\check-config.txt\" 2>&1 & type \"$guest_root\\check-config.txt\" & \"$guest_root\\rvbox.exe\" --service --config \"$guest_root\\client.toml\" > \"$guest_root\\direct-service-probe.txt\" 2>&1 & type \"$guest_root\\direct-service-probe.txt\" & if exist \"C:/ProgramData/RVBox/service-startup.log\" type \"C:/ProgramData/RVBox/service-startup.log\" & wevtutil qe System /q:\"*[System[(EventID=7000 or EventID=7009 or EventID=7031 or EventID=7034)]]\" /c:3 /rd:true /f:text & if exist \"C:/ProgramData/RVBox/test-logs/rvbox.log\" type \"C:/ProgramData/RVBox/test-logs/rvbox.log\" & echo RVBOX_GUEST_OK"
/d /s /c "sc.exe queryex RVBoxClient & sc.exe qc RVBoxClient & reg.exe query \"HKLM\\SYSTEM\\CurrentControlSet\\Services\\RVBoxClient\" /v ImagePath & reg.exe query \"HKLM\\SYSTEM\\CurrentControlSet\\Services\\RVBoxClient\" /v ObjectName & dir \"C:/ProgramData/RVBox\" & icacls \"C:/ProgramData/RVBox\" & certutil -hashfile \"$guest_root\\rvbox.exe\" SHA256 & \"$guest_root\\rvbox.exe\" --check-config --config \"$guest_root\\client.toml\" > \"$guest_root\\check-config.txt\" 2>&1 & type \"$guest_root\\check-config.txt\" & if exist \"C:/ProgramData/RVBox/service-startup.log\" type \"C:/ProgramData/RVBox/service-startup.log\" & wevtutil qe System /q:\"*[System[(EventID=7000 or EventID=7009 or EventID=7031 or EventID=7034)]]\" /c:3 /rd:true /f:text & if exist \"C:/ProgramData/RVBox/test-logs/rvbox.log\" type \"C:/ProgramData/RVBox/test-logs/rvbox.log\" & echo RVBOX_GUEST_OK"
;;
logs)
assert_identity
require_lease
[ "$(state)" = running ] || fail "logs requires a running prepared VM"
provisioner_run --exe 'C:\Windows\System32\cmd.exe' --wait-stdout --wait-stderr --unquoted-args -- \
/d /s /c "dir \"C:\\ProgramData\\RVBox\" & dir \"C:\\ProgramData\\RVBox\\test-logs\" & type \"C:\\ProgramData\\RVBox\\service-startup.log\" & type \"C:\\ProgramData\\RVBox\\test-logs\\rvbox.log\" & echo RVBOX_GUEST_OK"
;;
stop)
assert_identity
@@ -540,7 +548,22 @@ case "$action" in
attempt=$((attempt + 1))
sleep 1
done
[ "$(state)" = poweroff ] || fail "guest did not power off before reset"
# This is the exact named disposable fixture, with the matching
# run lease. Reset is its isolation boundary, not a graceful-stop
# diagnostic command: after a bounded ACPI attempt, discard only
# this guest's state so snapshot restore cannot strand the lane.
if [ "$(state)" != poweroff ]; then
printf 'reset: ACPI shutdown timed out; forcing disposable VM poweroff\n' >&2
VBoxManage controlvm "$vm" poweroff >/dev/null
attempt=0
while [ "$attempt" -lt 30 ]; do
[ "$(state)" = poweroff ] && break
attempt=$((attempt + 1))
sleep 1
done
[ "$(state)" = poweroff ] || fail "guest did not power off after forced reset"
printf 'reset-force-poweroff vm=%s\n' "$vm"
fi
fi
VBoxManage snapshot "$vm" restore "$snapshot" >/dev/null
assert_identity
@@ -580,7 +603,9 @@ accelerated_stage() {
printf 'stage: accelerated publish directory unavailable\n' >&2
return 1
}
stage_http_dir=$(mktemp -d "${TMPDIR:-/tmp}/rvbox-http-stage.XXXXXX")
# Keep all disposable test bytes inside the owning run directory; never
# consume global /tmp, which may belong to a different test or user.
stage_http_dir=$(mktemp -d "$run_root/.accelerated-stage.XXXXXX")
stage_http_xz=$stage_http_dir/rvbox.exe.xz
xz -T0 -3 -c "$bundle/rvbox.exe" >"$stage_http_xz"
stage_http_hash=$(sha256sum "$stage_http_xz" | awk '{print $1}')
@@ -615,7 +640,7 @@ copy_stage_file() {
target_name=$2
copy_attempt=1
while [ "$copy_attempt" -le 4 ]; do
if scp -q "$source_file" "$RVBOX_TEST_VBOX_HOST:$host_stage/$target_name"; then
if scp -q -o BatchMode=yes -o ConnectTimeout=10 -o ServerAliveInterval=10 -o ServerAliveCountMax=2 "$source_file" "$RVBOX_TEST_VBOX_HOST:$host_stage/$target_name"; then
return 0
fi
copy_attempt=$((copy_attempt + 1))
+35
View File
@@ -173,6 +173,17 @@ layer = "unit"
status = "implemented"
tests = ["internal/agentproto/validate_test.go:TestOutputChunkBoundedDecompression_HP_PROTO_05"]
[[requirements]]
id = "SEC-PROTO-01"
layer = "unit"
status = "implemented"
tests = [
"internal/agentproto/validate_test.go:FuzzDecodeEnvelopeBounded_SEC_PROTO_01",
"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",
]
[[requirements]]
id = "HP-PROTO-09"
layer = "unit"
@@ -608,6 +619,30 @@ layer = "unit"
status = "implemented"
tests = ["test/windowsnative/native_fixture_test.go:TestProductionComposeAssets_HP_OPS_01"]
[[requirements]]
id = "HP-OPS-02"
layer = "unit"
status = "implemented"
tests = ["internal/observability/rotate_test.go:TestRotatingFileBoundsArchives_HP_OPS_02"]
[[requirements]]
id = "BH-OPS-03"
layer = "unit"
status = "implemented"
tests = ["internal/observability/rotate_test.go:TestRotatingFileRejectsPartialRotationConfig_BH_OPS_03"]
[[requirements]]
id = "HP-OPS-04"
layer = "unit"
status = "implemented"
tests = ["internal/observability/rotate_test.go:TestFormatLogWritesStructuredBoundedRecords_HP_OPS_04"]
[[requirements]]
id = "HP-DISPATCH-10"
layer = "unit"
status = "implemented"
tests = ["internal/server/store/command_test.go:TestPendingStdinWaitsForDurableDispatch_HP_DISPATCH_06"]
[[requirements]]
id = "HP-STORE-01"
layer = "integration"
+1
View File
@@ -32,6 +32,7 @@ services:
- "${RVBOX_NATIVE_RUNTIME_DIR}/pki:/pki"
- ./certgen.sh:/fixture/certgen.sh:ro
entrypoint: ["/bin/sh", "/fixture/certgen.sh"]
networks: [native]
networks:
native:
labels: { rvbox.native.run_id: "${RVBOX_NATIVE_RUN_ID}" }
+18 -3
View File
@@ -29,6 +29,10 @@ func TestNativeFixtureAssets_HP_HARNESS_20(t *testing.T) {
"ACTIVE_USER_ELEVATED,ACTIVE_SYSTEM",
"assert_context local-service no local-service",
"assert_stdin_close",
"assert_signal_term",
"rvc kill TERM \"$client_id\" \"$issue\"",
"assert_server_restart_reconnect",
"compose restart server",
"rvc append \"$client_id\" \"$issue\" RVBOX_NATIVE_INPUT",
"rvc close-stdin \"$client_id\" \"$issue\"",
"clean --purge --yes",
@@ -40,13 +44,13 @@ func TestNativeFixtureAssets_HP_HARNESS_20(t *testing.T) {
}
}
testHost := read("scripts/windows/test-host")
for _, required := range []string{"xz -T0 -3", "accelerated_stage", "copy_stage_file", "--proxy", "--anyauth", "--continue-at", "rvbox.exe.xz", "retry_limit=4", "ConnectTimeout=10", "bundle transfer did not reach the expected SHA-256 manifest", "verified transfer_sha256"} {
for _, required := range []string{"xz -T0 -3", "accelerated_stage", "copy_stage_file", "--proxy", "--anyauth", "--continue-at", "rvbox.exe.xz", "retry_limit=4", "ConnectTimeout=10", "bundle transfer did not reach the expected SHA-256 manifest", "verified transfer_sha256", "reset-force-poweroff", "$run_root/.accelerated-stage.XXXXXX", "service-startup.log", "test-logs/rvbox.log"} {
if !strings.Contains(testHost, required) {
t.Fatalf("native test-host is missing compressed transfer contract %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"} {
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) {
t.Fatalf("native Compose fixture is missing %q", required)
}
@@ -56,6 +60,16 @@ func TestNativeFixtureAssets_HP_HARNESS_20(t *testing.T) {
if !strings.Contains(defaults, "//go:build !rvbox_native_test") || !strings.Contains(fixture, "//go:build rvbox_native_test") {
t.Fatal("fixture-only context faults are not separated from release builds")
}
service := read("cmd/rvbox/service_windows.go")
if !strings.Contains(service, "applyServiceDiagnosticsACL") || !strings.Contains(service, "D:P(A;;FA;;;SY)(A;;FA;;;BA)") {
t.Fatal("Windows service diagnostics are not readable by administrators")
}
testing := read("docs/testing.md")
for _, required := range []string{"server-process\nrestart", "new reconciled WSS session", "`logs` is the narrow read-only"} {
if !strings.Contains(testing, required) {
t.Fatalf("native workflow documentation is missing %q", required)
}
}
}
// TestProductionComposeAssets_HP_OPS_01 guards the deployment properties that
@@ -77,6 +91,7 @@ func TestProductionComposeAssets_HP_OPS_01(t *testing.T) {
"RVBOX_SERVER_CONFIG:-./server.toml",
"RVBOX_HTTPS_BIND:-0.0.0.0",
"nginx@sha256:65645c7bb6a0661892a8b03b89d0743208a18dd2f3f17a54ef4b76fb8e2f2a10",
"user: \"65532:65532\"",
"nofile:", "soft: 65536", "hard: 65536",
"condition: service_completed_successfully",
"condition: service_healthy",
@@ -87,7 +102,7 @@ func TestProductionComposeAssets_HP_OPS_01(t *testing.T) {
}
}
config := read("deploy/production/server.toml.example")
for _, required := range []string{"agent_listen = \"0.0.0.0:6899\"", "listen = \"0.0.0.0:6901\"", "enabled = false"} {
for _, required := range []string{"agent_listen = \"0.0.0.0:6899\"", "listen = \"0.0.0.0:6901\"", "enabled = false", "log_file = \"/var/lib/rvbox-server/logs/rvbox-server.log\"", "log_max_bytes = 10485760", "log_max_files = 5"} {
if !strings.Contains(config, required) {
t.Fatalf("production server example is missing %q", required)
}