Compare commits

...
6 Commits
30 changed files with 882 additions and 83 deletions
+11 -4
View File
@@ -6,6 +6,7 @@ import (
"errors" "errors"
"flag" "flag"
"fmt" "fmt"
"io"
"log" "log"
"net" "net"
"net/http" "net/http"
@@ -58,6 +59,12 @@ func run(configPath string) error {
if err != nil { if err != nil {
return err 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{ persistence, err := store.Open(context.Background(), store.Options{
DataDir: configured.Server.DataDir, BusyTimeout: configured.Storage.SQLiteBusyTimeout, DataDir: configured.Server.DataDir, BusyTimeout: configured.Storage.SQLiteBusyTimeout,
SegmentTargetSize: configured.Storage.SegmentTargetBytes, SegmentTargetSize: configured.Storage.SegmentTargetBytes,
@@ -84,7 +91,7 @@ func run(configPath string) error {
if configured.Observability.Listen != "" { if configured.Observability.Listen != "" {
healthListener, err = net.Listen("tcp", configured.Observability.Listen) healthListener, err = net.Listen("tcp", configured.Observability.Listen)
if err != nil { 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 { if healthListener != nil {
@@ -104,10 +111,10 @@ func run(configPath string) error {
Scope: store.IncidentScopeGlobal, ScopeKey: "server-startup-recovery", Summary: "startup storage recovery failed", Scope: store.IncidentScopeGlobal, ScopeKey: "server-startup-recovery", Summary: "startup storage recovery failed",
Evidence: []byte(recoverErr.Error()), AutomaticallyRepairable: false, Evidence: []byte(recoverErr.Error()), AutomaticallyRepairable: false,
}); incidentErr != nil { }); 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 return
} }
health.SetReady(true) health.SetReady(true)
@@ -143,7 +150,7 @@ func run(configPath string) error {
} }
defer rpcListener.Close() defer rpcListener.Close()
if configured.JSONRPC.NonLoopbackBind { 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} 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 { if err != nil {
return err 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() health := observability.New()
go func() { 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 { 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() }, Jitter: agent.CryptoJitter, Now: func() time.Time { return time.Now().UTC() },
OnDispatch: executor.Dispatch, OnScriptReady: executor.ScriptReady, OnStdin: executor.Stdin, OnDispatch: executor.Dispatch, OnScriptReady: executor.ScriptReady, OnStdin: executor.Stdin,
OnCloseStdin: executor.CloseStdin, OnSignal: executor.Signal, OnTerminate: executor.Terminate, EventReady: eventReady, 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 { }); runErr != nil && ctx.Err() == nil && diagnostics != nil {
_, _ = fmt.Fprintf(diagnostics, "rvbox client session stopped: %v\n", runErr) _, _ = 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" clientwindows "github.com/rvbox/rvbox/internal/client/supervisor/windows"
"github.com/rvbox/rvbox/internal/client/windowsservice" "github.com/rvbox/rvbox/internal/client/windowsservice"
"github.com/rvbox/rvbox/internal/client/windowstray" "github.com/rvbox/rvbox/internal/client/windowstray"
"golang.org/x/sys/windows"
"golang.org/x/sys/windows/svc" "golang.org/x/sys/windows/svc"
) )
@@ -90,6 +91,16 @@ func openServiceDiagnostics(configPath string, diagnostics io.Writer) (io.Writer
} }
return diagnostics, func() {} 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 { if diagnostics == nil {
return file, func() { _ = file.Close() } 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() } 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 { func runTray(configPath string, diagnostics io.Writer) error {
_ = configPath // the tray obtains the canonical paths from the service. _ = configPath // the tray obtains the canonical paths from the service.
return windowstray.Run(context.Background(), diagnostics) return windowstray.Run(context.Background(), diagnostics)
+3
View File
@@ -27,6 +27,9 @@ services:
server: server:
image: "${RVBOX_SERVER_IMAGE:?set RVBOX_SERVER_IMAGE to a pinned rvbox-server image}" 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 restart: unless-stopped
command: ["--config", "/etc/rvbox/server.toml"] command: ["--config", "/etc/rvbox/server.toml"]
read_only: true read_only: true
+6
View File
@@ -19,3 +19,9 @@ enabled = false
[observability] [observability]
# nginx exposes only /livez and /readyz, not the metrics listener itself. # nginx exposes only /livez and /readyz, not the metrics listener itself.
listen = "0.0.0.0:6901" 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 4. The server queues work while a client is offline and dispatches it only when
the active session advertises capacity. A replacement session immediately the active session advertises capacity. A replacement session immediately
performs bidirectional reconciliation between the server's non-terminal set performs bidirectional reconciliation between the server's non-terminal set
and the client's complete retained-command set. The server returns explicit and the client's complete retained-command set. That comparison uses one
local terminate/discard decisions before new dispatch begins. 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 ## 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 result before dispatch. A client with unresolved essential-store corruption
cannot complete reconciliation or accept work. 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: Implement reconciliation as this explicit matrix:
| Server state | Client evidence | Durable result | | 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 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: Build all supported binaries without installing `make` or Go on the host:
```sh ```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 containers. The self-signed server certificate is intentionally accepted by
the v1 client without a test CA. It then drives the installed SCM service 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 through the server's real Unix control socket and verifies every Windows
execution context. The tagged binary's controlled pre-launch failures are execution context, ordered stdin close, TERM delivery, and a server-process
limited to the test fixture; a release binary rejects that switch. 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 Successful runs collect bounded artifacts, remove only their labeled Compose
project, and restore the exact clean snapshot. A failed or --keep run stays 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. log, or artifact.
The native lifecycle is `status`, `prepare`, `stage`, `install`, `run`, 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 restores the clean baseline, starts headless, waits for Guest Additions, and
proves that `RVBoxClient` is absent. `stage` copies a versioned non-secret test 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. 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. the real `rvbox.exe --install-service` path and proves completion through SCM.
`run` is for reconfiguration/restart scenarios after that first installation. `run` is for reconfiguration/restart scenarios after that first installation.
Neither action invokes the GUI-subsystem executable directly with the normal Neither action invokes the GUI-subsystem executable directly with the normal
Guest Control account. `collect` obtains only bounded/redacted artifacts, and Guest Control account. `collect` obtains only bounded/redacted artifacts.
`reset` restores the exact clean baseline and leaves the VM powered off. `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 This service-driven protocol is required because VirtualBox Guest Control
7.2.16 does not reliably complete a direct GUI-subsystem `rvbox.exe` run; 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) { func TestEnvelopeSessionFencingShape_BH_SES_01(t *testing.T) {
t.Parallel() t.Parallel()
+29 -22
View File
@@ -35,6 +35,10 @@ type RunnerOptions struct {
OnSignal func(context.Context, Session, *rvboxv1.SignalCommand) error OnSignal func(context.Context, Session, *rvboxv1.SignalCommand) error
OnScriptReady func(context.Context, Session, domain.UUID) error OnScriptReady func(context.Context, Session, domain.UUID) error
OnTerminate func(context.Context, 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 // EventReady wakes the active session after a supervisor worker appends a
// durable event. The network loop remains the sole writer; a reconnect can // durable event. The network loop remains the sole writer; a reconnect can
// safely ignore a stale notification because replay reads the spool again. // 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 return nil
} }
sessionStarted := options.Now() sessionStarted := options.Now()
_ = runOnce(ctx, options) sessionErr := runOnce(ctx, options)
if ctx.Err() != nil { if ctx.Err() != nil {
return nil return nil
} }
if sessionErr != nil && options.OnSessionError != nil {
options.OnSessionError(sessionErr)
}
if options.Now().Sub(sessionStarted) >= options.Backoff.StableReset { if options.Now().Sub(sessionStarted) >= options.Backoff.StableReset {
failures = 0 failures = 0
} }
@@ -152,7 +159,7 @@ func RunOnce(ctx context.Context, options RunnerOptions) error {
func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) { func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
transport, err := options.Dial(ctx) transport, err := options.Dial(ctx)
if err != nil { if err != nil {
return err return fmt.Errorf("dial agent server: %w", err)
} }
defer func() { defer func() {
if closeErr := transport.Close(); resultErr == nil && closeErr != nil { 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) session, err := Handshake(ctx, transport, options.Hello, limits)
if err != nil { if err != nil {
return err return fmt.Errorf("agent handshake: %w", err)
} }
snapshot, err := options.Store.ReconcileSnapshot(ctx) snapshot, err := options.Store.ReconcileSnapshot(ctx)
if err != nil { if err != nil {
@@ -174,7 +181,7 @@ func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
} }
result, err := Reconcile(ctx, transport, session, snapshot, limits) result, err := Reconcile(ctx, transport, session, snapshot, limits)
if err != nil { if err != nil {
return err return fmt.Errorf("exchange reconciliation: %w", err)
} }
terminated, err := ApplyReconcileResult(ctx, options.Store, result, options.Now()) terminated, err := ApplyReconcileResult(ctx, options.Store, result, options.Now())
if err != nil { if err != nil {
@@ -183,16 +190,16 @@ func runOnce(ctx context.Context, options RunnerOptions) (resultErr error) {
for _, issue := range terminated { for _, issue := range terminated {
if options.OnTerminate != nil { if options.OnTerminate != nil {
if err := options.OnTerminate(ctx, issue); err != 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) sentEvents := make(map[domain.UUID]uint64)
if err := replayEvents(ctx, transport, options.Store, session, snapshot, limits, sentEvents); err != nil { 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 { 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) return serveActive(ctx, transport, options, session, sentEvents)
} }
@@ -323,20 +330,20 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
case <-ctx.Done(): case <-ctx.Done():
return nil return nil
case err := <-readErrors: case err := <-readErrors:
return err return fmt.Errorf("read active agent frame: %w", err)
case issue := <-options.EventReady: case issue := <-options.EventReady:
if issue == (domain.UUID{}) { if issue == (domain.UUID{}) {
continue continue
} }
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != nil { 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 continue
case encoded = <-frames: case encoded = <-frames:
} }
envelope, err := agentproto.DecodeEnvelope(encoded, limits, rvboxv1.Platform_PLATFORM_WINDOWS) envelope, err := agentproto.DecodeEnvelope(encoded, limits, rvboxv1.Platform_PLATFORM_WINDOWS)
if err != nil { if err != nil {
return err return fmt.Errorf("decode active agent frame: %w", err)
} }
if envelope.GetSessionId() != session.ID || envelope.GetSessionGeneration() != session.Generation { if envelope.GetSessionId() != session.ID || envelope.GetSessionGeneration() != session.Generation {
return ErrProtocolHandshake return ErrProtocolHandshake
@@ -344,25 +351,25 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
switch { switch {
case envelope.GetCommandDispatch() != nil: case envelope.GetCommandDispatch() != nil:
if err := handleDispatch(ctx, transport, options, session, envelope.GetCommandDispatch(), limits); err != 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()) issue, parseErr := domain.ParseUUIDv7(envelope.GetCommandDispatch().GetIssueUuid())
if parseErr == nil { if parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != 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: case envelope.GetEventAck() != nil:
if err := ApplyEventAck(ctx, options.Store, envelope.GetEventAck()); err != 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: case envelope.GetScriptChunk() != nil:
if err := handleScriptChunk(ctx, transport, options.Store, session, envelope.GetScriptChunk(), limits, options.Now, sent); err != 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: case envelope.GetScriptCommit() != nil:
if err := handleScriptCommit(ctx, transport, options.Store, session, envelope.GetScriptCommit(), limits, options.Now, sent); err != 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 { if options.OnScriptReady != nil {
issue, err := domain.ParseUUIDv7(envelope.GetScriptCommit().GetIssueUuid()) issue, err := domain.ParseUUIDv7(envelope.GetScriptCommit().GetIssueUuid())
@@ -370,40 +377,40 @@ func serveActive(ctx context.Context, transport Transport, options RunnerOptions
return err return err
} }
if err := options.OnScriptReady(ctx, session, issue); err != nil { 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: case envelope.GetStdinWrite() != nil:
if options.OnStdin != nil { if options.OnStdin != nil {
if err := options.OnStdin(ctx, session, envelope.GetStdinWrite()); err != 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 issue, parseErr := domain.ParseUUIDv7(envelope.GetStdinWrite().GetIssueUuid()); parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != 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: case envelope.GetCloseStdin() != nil:
if options.OnCloseStdin != nil { if options.OnCloseStdin != nil {
if err := options.OnCloseStdin(ctx, session, envelope.GetCloseStdin()); err != 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 issue, parseErr := domain.ParseUUIDv7(envelope.GetCloseStdin().GetIssueUuid()); parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != 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: case envelope.GetSignalCommand() != nil:
if options.OnSignal != nil { if options.OnSignal != nil {
if err := options.OnSignal(ctx, session, envelope.GetSignalCommand()); err != 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 issue, parseErr := domain.ParseUUIDv7(envelope.GetSignalCommand().GetIssueUuid()); parseErr == nil {
if err := flushEvents(ctx, transport, options.Store, session, issue, limits, sent); err != 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: 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) { func TestFlushEventsAssignsAndSendsOnlyUnacknowledgedRows_HP_RUNTIME_03(t *testing.T) {
ctx := context.Background() ctx := context.Background()
store, err := spool.Open(ctx, spool.Options{DataDir: filepath.Join(t.TempDir(), "spool"), BusyTimeout: time.Second}) 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{ Observability: observabilityFile{
Listen: "127.0.0.1:6901", LivenessPath: "/livez", ReadinessPath: "/readyz", 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") t.Fatal("effective server configuration is nondeterministic")
} }
digest := sha256.Sum256(first) digest := sha256.Sum256(first)
const goldenSHA256 = "55c5fec2c10bc27372353a60f317c5a2959f8087685cf7c11202bb12dc59aa28" const goldenSHA256 = "65c938efd935ce0c0d2354cdb62c38a2dc177755520d1ae1e5951902a70dd323"
if got := hex.EncodeToString(digest[:]); got != goldenSHA256 { if got := hex.EncodeToString(digest[:]); got != goldenSHA256 {
t.Fatalf("effective server configuration golden changed: got %s", got) 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) handler.writeRPCError(response, nil, jsonRPCParseError, "could not read request", nil)
return return
} }
if int64(len(body)) > handler.MaxBody { envelope, errorCode, errorMessage := decodeJSONRPCRequest(body, handler.MaxBody)
handler.writeRPCError(response, nil, jsonRPCInvalid, "request body exceeds limit", nil) if errorCode != 0 {
return handler.writeRPCError(response, nil, errorCode, errorMessage, nil)
}
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)
return return
} }
result, callErr := handler.call(request.Context(), envelope.Method, envelope.Params) result, callErr := handler.call(request.Context(), envelope.Method, envelope.Params)
@@ -87,6 +76,29 @@ type jsonRPCRequest struct {
Params json.RawMessage `json:"params"` 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 { type jsonRPCResponse struct {
JSONRPC string `json:"jsonrpc"` JSONRPC string `json:"jsonrpc"`
ID json.RawMessage `json:"id"` 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) 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 return
} }
hello := helloEnvelope.GetClientHello() 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()) instanceID, err := domain.ParseUUIDv7(hello.GetClientInstanceId())
if err != nil { if err != nil {
server.close(connection, websocket.StatusPolicyViolation, "invalid client instance ID") 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") server.close(connection, websocket.StatusPolicyViolation, "invalid initial client capacity")
return 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) encodedSessionID := encodeSessionID(sessionID)
welcome, err := proto.Marshal(&rvboxv1.AgentEnvelope{ welcome, err := proto.Marshal(&rvboxv1.AgentEnvelope{
SessionId: encodedSessionID, SessionGeneration: registration.Generation, 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") server.close(connection, websocket.StatusInternalError, "could not queue welcome")
return return
} }
targets, err := server.Store.ReconcileTargets(parent, hello.GetClientId()) targets, err := server.Store.ReconcileTargetsAt(parent, hello.GetClientId(), reconcileBoundary)
if err != nil { if err != nil {
server.close(connection, websocket.StatusInternalError, "could not build reconciliation request") server.close(connection, websocket.StatusInternalError, "could not build reconciliation request")
return return
@@ -237,7 +241,7 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
return return
} }
if snapshot := envelope.GetReconcileSnapshot(); snapshot != nil { 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 { if reconcileErr != nil {
server.close(connection, websocket.StatusPolicyViolation, "reconciliation failed") server.close(connection, websocket.StatusPolicyViolation, "reconciliation failed")
return return
@@ -409,7 +413,7 @@ func eventType(event *rvboxv1.CommandEvent) uint16 {
// enqueueNextDispatch records the queued-to-dispatched transition before // enqueueNextDispatch records the queued-to-dispatched transition before
// exposing work to the network. A full data lane is a pre-write failure, so // 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. // 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()) candidate, err := server.Store.ClaimNextDispatch(ctx, clientID, generation, server.now())
if err != nil || candidate == nil { if err != nil || candidate == nil {
return domain.UUID{}, false, err return domain.UUID{}, false, err
@@ -447,7 +451,11 @@ func (server *AgentServer) enqueueNextDispatch(ctx context.Context, queue *Write
if beforeEnqueue != nil { if beforeEnqueue != nil {
beforeEnqueue(candidate) 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) sent, requeueErr := requeue(ErrDispatchDataFull)
return candidate.IssueUUID, sent, requeueErr 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 // lane. Session-local sent tracking suppresses duplicate frames while a live
// connection remains usable; reconnecting naturally replays unacknowledged // connection remains usable; reconnecting naturally replays unacknowledged
// writes from storage. // 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) intents, err := server.Store.PendingStdin(ctx, clientID)
if err != nil { if err != nil {
return false, err return false, err
} }
for _, intent := range intents { for _, intent := range intents {
sentMu.Lock()
written, waitingForDispatch := dispatchWritten[intent.IssueUUID.String()]
sentMu.Unlock()
if waitingForDispatch && !written {
continue
}
key := stdinIntentKey(intent.IssueUUID, intent.WriteSeq) key := stdinIntentKey(intent.IssueUUID, intent.WriteSeq)
sentMu.Lock() sentMu.Lock()
_, alreadySent := sent[key] _, 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) 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) intents, err := server.Store.PendingSignals(ctx, clientID, generation)
if err != nil { if err != nil {
return false, err return false, err
} }
for _, intent := range intents { 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) key := signalIntentKey(intent.IssueUUID, intent.CommandRevision, intent.Signal)
sentMu.Lock() sentMu.Lock()
_, alreadySent := sent[key] _, 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 // dispatchLoop is the per-session serialized dispatcher. It waits for a
// complete reconciliation result before consuming queued work, then coalesces // complete reconciliation result before consuming queued work, then coalesces
// wakeups from local control RPCs, capacity advertisements, and acceptances. // 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 { select {
case <-reconciled: case <-reconciled:
case <-ctx.Done(): case <-ctx.Done():
@@ -621,7 +645,7 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue,
} }
for { for {
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 { if stdinErr != nil {
server.closeForDispatchFailure(cancel) server.closeForDispatchFailure(cancel)
return return
@@ -629,7 +653,7 @@ func (server *AgentServer) dispatchLoop(ctx context.Context, queue *WriterQueue,
if stdinQueued { if stdinQueued {
continue 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 { if signalErr != nil {
server.closeForDispatchFailure(cancel) server.closeForDispatchFailure(cancel)
return 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) { issue, sent, err := server.enqueueNextDispatch(ctx, queue, clientID, platform, sessionID, generation, func(candidate *store.DispatchCandidate) {
capacityMu.Lock() capacityMu.Lock()
reservations[candidate.IssueUUID.String()] = lane reservations[candidate.IssueUUID.String()] = lane
dispatchWritten[candidate.IssueUUID.String()] = false
if candidate.ScriptPresent { if candidate.ScriptPresent {
scriptTransfers[candidate.IssueUUID.String()] = &scriptTransfer{Body: append([]byte(nil), candidate.ScriptContent...), Digest: sha256.Sum256(candidate.ScriptContent)} scriptTransfers[candidate.IssueUUID.String()] = &scriptTransfer{Body: append([]byte(nil), candidate.ScriptContent...), Digest: sha256.Sum256(candidate.ScriptContent)}
} }
inserted = true inserted = true
capacityMu.Unlock() capacityMu.Unlock()
}, func(issue domain.UUID) {
capacityMu.Lock()
dispatchWritten[issue.String()] = true
capacityMu.Unlock()
if signalWake != nil {
signalWake()
}
}) })
if !sent { if !sent {
capacityMu.Lock() capacityMu.Lock()
if inserted { if inserted {
delete(reservations, issue.String()) delete(reservations, issue.String())
delete(dispatchWritten, issue.String())
delete(scriptTransfers, issue.String()) delete(scriptTransfers, issue.String())
} }
capacity.Release(lane) capacity.Release(lane)
@@ -706,6 +739,9 @@ func (server *AgentServer) writeLoop(ctx context.Context, connection *websocket.
if frame.Written != nil { if frame.Written != nil {
close(frame.Written) close(frame.Written)
} }
if frame.OnWritten != nil {
frame.OnWritten()
}
continue continue
} }
if !errors.Is(err, context.DeadlineExceeded) { 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) { func TestAgentServerRejectsHostileWireInputs_BH_SES_04(t *testing.T) {
server, cleanup := newTestAgentServer(t) server, cleanup := newTestAgentServer(t)
defer cleanup() defer cleanup()
@@ -174,6 +225,11 @@ func TestWireEventAppendCarriesClientBinding_HP_EVENT_01(t *testing.T) {
} }
func newTestAgentServer(t *testing.T) (*httptest.Server, func()) { 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() t.Helper()
dataDirectory := filepath.Join(t.TempDir(), "store") dataDirectory := filepath.Join(t.TempDir(), "store")
if err := os.Mkdir(dataDirectory, 0o700); err != nil { 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, HeartbeatIdle: time.Second, LivenessTimeout: 2 * time.Second,
} }
httpServer := httptest.NewServer(agent) httpServer := httptest.NewServer(agent)
return httpServer, func() { return httpServer, agent, func() {
httpServer.Close() httpServer.Close()
if err := persistence.Close(); err != nil { if err := persistence.Close(); err != nil {
t.Errorf("close persistence: %v", err) 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. // Written is closed by the sole socket writer after a successful write.
// It is used only for protocol barriers such as reconciliation-before-work. // It is used only for protocol barriers such as reconciliation-before-work.
Written chan<- struct{} 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 // 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) { func TestLateDispatchAcceptanceAndEventRetainExpiryContradiction_BH_DISPATCH_09(t *testing.T) {
t.Parallel() t.Parallel()
ctx := context.Background() 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 // 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. // UUIDs and durable cursors must be compared before fresh dispatch is enabled.
func (store *Store) ReconcileTargets(ctx context.Context, clientID string) ([]*rvboxv1.ReconcileTarget, error) { 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 == "" { if clientID == "" {
return nil, errors.New("client ID is required") return nil, errors.New("client ID is required")
} }
rows, err := store.db.QueryContext(ctx, `SELECT issue_uuid, last_event_seq, revision, immutable_request_sha256 query := `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) 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 { if err != nil {
return nil, err 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 // incidented. Retained non-terminal rows are retargeted to this generation so
// late events from the previous connection cannot advance the command. // 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) { 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 { if generation == 0 {
return nil, errors.New("session generation is required") 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 { 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) { 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 == "" { if clientID == "" {
return nil, errors.New("client ID is required") return nil, errors.New("client ID is required")
} }
@@ -112,7 +140,7 @@ func (store *Store) reconcileClientSnapshot(ctx context.Context, clientID string
return nil, err return nil, err
} }
defer tx.Rollback() defer tx.Rollback()
serverRows, tombstones, err := loadReconcileRows(ctx, tx, clientID) serverRows, tombstones, err := loadReconcileRows(ctx, tx, clientID, receiptBoundary)
if err != nil { if err != nil {
store.writeMu.Unlock() store.writeMu.Unlock()
return nil, err return nil, err
@@ -215,8 +243,14 @@ type reconcileIncident struct {
dataLoss bool dataLoss bool
} }
func loadReconcileRows(ctx context.Context, tx *sql.Tx, clientID string) (map[string]reconcileServerRow, map[string][]byte, error) { func loadReconcileRows(ctx context.Context, tx *sql.Tx, clientID string, receiptBoundary time.Time) (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) 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 { if err != nil {
return nil, nil, err 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 { func mustReconcileIssue(t *testing.T, value string) domain.UUID {
t.Helper() t.Helper()
issue, err := domain.ParseUUIDv7(value) issue, err := domain.ParseUUIDv7(value)
+7 -5
View File
@@ -46,10 +46,12 @@ type StdinIntent struct {
Close bool Close bool
} }
// PendingStdin returns unacknowledged input intents for one client. Delivery is // PendingStdin returns unacknowledged input intents only after the command has
// intentionally tracked by the session dispatcher, not by this durable query: // entered DISPATCHED state. This durable gate prevents append/close controls
// a lost connection simply causes the next session to replay the same write // from overtaking their initial CommandDispatch on a fresh session. Delivery is
// sequence, which the client acknowledges idempotently. // 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) { func (store *Store) PendingStdin(ctx context.Context, clientID string) ([]StdinIntent, error) {
if clientID == "" { if clientID == "" {
return nil, errors.New("client ID is required") 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 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 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) ORDER BY c.issue_time, c.issue_uuid, w.write_seq`, clientID)
if err != nil { if err != nil {
return nil, err 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" "$@" 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() { rvc() {
compose exec -T server /opt/rvbox/rvc --socket /run/rvbox/server.sock "$@" 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" 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() { collect() {
if [ "$vm_prepared" = yes ]; then if [ "$vm_prepared" = yes ]; then
"$repo_root/scripts/windows/test-host" collect --run-id "$run_id" || true "$repo_root/scripts/windows/test-host" collect --run-id "$run_id" || true
@@ -190,6 +262,7 @@ collect() {
clean() { clean() {
compose down --volumes --remove-orphans || true compose down --volumes --remove-orphans || true
remove_project_networks
# A reset is the isolation boundary for the next run. Do not conceal a # A reset is the isolation boundary for the next run. Do not conceal a
# failed shutdown/snapshot restore behind a successful-looking `clean`: # failed shutdown/snapshot restore behind a successful-looking `clean`:
# callers must repair or explicitly inspect the retained VM lease first. # 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" "$repo_root/scripts/windows/test-host" run --run-id "$run_id" --endpoint "$endpoint_host:$port"
wait_client wait_client
assert_stdin_close assert_stdin_close
assert_signal_term
assert_server_restart_reconnect
assert_context active-user no active-user assert_context active-user no active-user
assert_context active-user-elevated yes active-user-elevated 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 "$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 logoff log off the sole active fixture user; use only after service installation
collect copy bounded guest artifacts to the local test-run directory 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 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 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 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 recover read-only fixture/run-state check for a stopped-resumable run
@@ -87,7 +88,7 @@ while [ "$#" -gt 0 ]; do
esac esac
done 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 if [ "$action" != status ]; then
[ -n "$run_id" ] || fail "$action requires --run-id" [ -n "$run_id" ] || fail "$action requires --run-id"
safe_id "$run_id" safe_id "$run_id"
@@ -169,7 +170,7 @@ remote() {
# retry. Lifecycle transitions remain single-attempt: their caller must # retry. Lifecycle transitions remain single-attempt: their caller must
# inspect/recover rather than risk a duplicate reset, shutdown, or logoff. # inspect/recover rather than risk a duplicate reset, shutdown, or logoff.
case $remote_action in 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 retry_limit=4
;; ;;
esac esac
@@ -488,9 +489,9 @@ case "$action" in
require_lease require_lease
install -d -m 700 "$host_stage/artifacts" install -d -m 700 "$host_stage/artifacts"
if [ "$(state)" = running ]; then if [ "$(state)" = running ]; then
guest_run --exe 'C:\Windows\System32\cmd.exe' --wait-stdout --wait-stderr --unquoted-args -- \ 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 & echo RVBOX_GUEST_OK" >/dev/null || true /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 "$guest_user" --passwordfile "$password_file" \ VBoxManage guestcontrol "$vm" --username "$provisioner_user" --passwordfile "$provisioner_password_file" \
copyfrom "$guest_root" "$host_stage/artifacts" --recursive </dev/null >/dev/null 2>&1 || true copyfrom "$guest_root" "$host_stage/artifacts" --recursive </dev/null >/dev/null 2>&1 || true
fi fi
printf 'collected host_stage=%s/artifacts\n' "$host_stage" printf 'collected host_stage=%s/artifacts\n' "$host_stage"
@@ -500,7 +501,14 @@ case "$action" in
require_lease require_lease
[ "$(state)" = running ] || fail "inspect requires a running prepared VM" [ "$(state)" = running ] || fail "inspect requires a running prepared VM"
guest_run --exe 'C:\Windows\System32\cmd.exe' --wait-stdout --wait-stderr --unquoted-args -- \ 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) stop)
assert_identity assert_identity
@@ -540,7 +548,22 @@ case "$action" in
attempt=$((attempt + 1)) attempt=$((attempt + 1))
sleep 1 sleep 1
done 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 fi
VBoxManage snapshot "$vm" restore "$snapshot" >/dev/null VBoxManage snapshot "$vm" restore "$snapshot" >/dev/null
assert_identity assert_identity
@@ -580,7 +603,9 @@ accelerated_stage() {
printf 'stage: accelerated publish directory unavailable\n' >&2 printf 'stage: accelerated publish directory unavailable\n' >&2
return 1 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 stage_http_xz=$stage_http_dir/rvbox.exe.xz
xz -T0 -3 -c "$bundle/rvbox.exe" >"$stage_http_xz" xz -T0 -3 -c "$bundle/rvbox.exe" >"$stage_http_xz"
stage_http_hash=$(sha256sum "$stage_http_xz" | awk '{print $1}') stage_http_hash=$(sha256sum "$stage_http_xz" | awk '{print $1}')
@@ -615,7 +640,7 @@ copy_stage_file() {
target_name=$2 target_name=$2
copy_attempt=1 copy_attempt=1
while [ "$copy_attempt" -le 4 ]; do 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 return 0
fi fi
copy_attempt=$((copy_attempt + 1)) copy_attempt=$((copy_attempt + 1))
+35
View File
@@ -173,6 +173,17 @@ layer = "unit"
status = "implemented" status = "implemented"
tests = ["internal/agentproto/validate_test.go:TestOutputChunkBoundedDecompression_HP_PROTO_05"] 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]] [[requirements]]
id = "HP-PROTO-09" id = "HP-PROTO-09"
layer = "unit" layer = "unit"
@@ -608,6 +619,30 @@ layer = "unit"
status = "implemented" status = "implemented"
tests = ["test/windowsnative/native_fixture_test.go:TestProductionComposeAssets_HP_OPS_01"] 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]] [[requirements]]
id = "HP-STORE-01" id = "HP-STORE-01"
layer = "integration" layer = "integration"
+1
View File
@@ -32,6 +32,7 @@ services:
- "${RVBOX_NATIVE_RUNTIME_DIR}/pki:/pki" - "${RVBOX_NATIVE_RUNTIME_DIR}/pki:/pki"
- ./certgen.sh:/fixture/certgen.sh:ro - ./certgen.sh:/fixture/certgen.sh:ro
entrypoint: ["/bin/sh", "/fixture/certgen.sh"] entrypoint: ["/bin/sh", "/fixture/certgen.sh"]
networks: [native]
networks: networks:
native: native:
labels: { rvbox.native.run_id: "${RVBOX_NATIVE_RUN_ID}" } 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", "ACTIVE_USER_ELEVATED,ACTIVE_SYSTEM",
"assert_context local-service no local-service", "assert_context local-service no local-service",
"assert_stdin_close", "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 append \"$client_id\" \"$issue\" RVBOX_NATIVE_INPUT",
"rvc close-stdin \"$client_id\" \"$issue\"", "rvc close-stdin \"$client_id\" \"$issue\"",
"clean --purge --yes", "clean --purge --yes",
@@ -40,13 +44,13 @@ func TestNativeFixtureAssets_HP_HARNESS_20(t *testing.T) {
} }
} }
testHost := read("scripts/windows/test-host") 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) { if !strings.Contains(testHost, required) {
t.Fatalf("native test-host is missing compressed transfer contract %q", required) t.Fatalf("native test-host is missing compressed transfer contract %q", required)
} }
} }
compose := read("test/linux-server/compose.yaml") 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) { if !strings.Contains(compose, required) {
t.Fatalf("native Compose fixture is missing %q", 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") { 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") 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 // 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_SERVER_CONFIG:-./server.toml",
"RVBOX_HTTPS_BIND:-0.0.0.0", "RVBOX_HTTPS_BIND:-0.0.0.0",
"nginx@sha256:65645c7bb6a0661892a8b03b89d0743208a18dd2f3f17a54ef4b76fb8e2f2a10", "nginx@sha256:65645c7bb6a0661892a8b03b89d0743208a18dd2f3f17a54ef4b76fb8e2f2a10",
"user: \"65532:65532\"",
"nofile:", "soft: 65536", "hard: 65536", "nofile:", "soft: 65536", "hard: 65536",
"condition: service_completed_successfully", "condition: service_completed_successfully",
"condition: service_healthy", "condition: service_healthy",
@@ -87,7 +102,7 @@ func TestProductionComposeAssets_HP_OPS_01(t *testing.T) {
} }
} }
config := read("deploy/production/server.toml.example") 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) { if !strings.Contains(config, required) {
t.Fatalf("production server example is missing %q", required) t.Fatalf("production server example is missing %q", required)
} }