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