Compare commits
6
Commits
d753e8b698
...
46185f1f6c
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
46185f1f6c | ||
|
|
684981c235 | ||
|
|
f6f900e597 | ||
|
|
418681d38e | ||
|
|
f86abcecb5 | ||
|
|
e79f882993 |
@@ -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}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|
||||||
|
|||||||
@@ -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
@@ -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;
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|
||||||
|
|||||||
@@ -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:
|
||||||
|
|||||||
@@ -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})
|
||||||
|
|||||||
@@ -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,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
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"`
|
||||||
|
|||||||
@@ -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)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
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" "$@"
|
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
|
||||||
|
|||||||
@@ -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))
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -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}" }
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user