diff --git a/cmd/rvbox-server/main.go b/cmd/rvbox-server/main.go index 4a9a8ab..a46b489 100644 --- a/cmd/rvbox-server/main.go +++ b/cmd/rvbox-server/main.go @@ -109,6 +109,7 @@ func run(configPath string) error { } recoveryContext, cancelRecovery := context.WithCancel(context.Background()) defer cancelRecovery() + go publishServerTelemetry(recoveryContext, persistence, health) // Recovery is deliberately asynchronous: liveness and incident inspection // remain available while committed-range checks run. Readiness becomes // true only after the real SQLite/segment recovery completes. @@ -239,3 +240,24 @@ func run(configPath string) error { return httpErr } } + +func publishServerTelemetry(ctx context.Context, persistence *store.Store, health *observability.Health) { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + for { + sampleContext, cancel := context.WithTimeout(ctx, time.Second) + telemetry, err := persistence.Telemetry(sampleContext) + cancel() + if err != nil { + health.Inc("telemetry_read_failure") + } else { + health.SetGauge("queue_depth", float64(telemetry.QueuedCommands)) + health.SetGauge("server_charged_bytes", float64(telemetry.ChargedBytes)) + } + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + } +} diff --git a/cmd/rvbox/main.go b/cmd/rvbox/main.go index 13770b0..a03e154 100644 --- a/cmd/rvbox/main.go +++ b/cmd/rvbox/main.go @@ -180,6 +180,7 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ return fmt.Errorf("open client durable state: %w", err) } defer state.Close() + go publishClientTelemetry(ctx, state, health) if diagnostics != nil { _, _ = fmt.Fprintf(diagnostics, "rvbox client service initialized for %s\n", configured.Client.ServerURL) } @@ -250,6 +251,27 @@ func runClientDaemon(ctx context.Context, configPath string, diagnostics io.Writ return nil } +func publishClientTelemetry(ctx context.Context, state *spool.Store, health *observability.Health) { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + for { + sampleContext, cancel := context.WithTimeout(ctx, time.Second) + telemetry, err := state.Telemetry(sampleContext) + cancel() + if err != nil { + health.Inc("telemetry_read_failure") + } else { + health.SetGauge("queue_depth", float64(telemetry.QueuedCommands)) + health.SetGauge("client_spool_bytes", float64(telemetry.ChargedBytes)) + } + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + } +} + func clientJobProfiles(profiles config.Profiles) map[string]clientwindows.JobProfile { return map[string]clientwindows.JobProfile{ rvboxv1.ExecutionProfile_EXECUTION_PROFILE_LIGHT.String(): toJobProfile(profiles.Light), diff --git a/internal/client/spool/telemetry.go b/internal/client/spool/telemetry.go new file mode 100644 index 0000000..ef51d97 --- /dev/null +++ b/internal/client/spool/telemetry.go @@ -0,0 +1,34 @@ +package spool + +import ( + "context" + "errors" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +// Telemetry is the bounded durable-spool state exposed as client gauges. It +// contains no command body, output, UUID, or client identity. +type Telemetry struct { + QueuedCommands uint64 + ChargedBytes uint64 +} + +func (store *Store) Telemetry(ctx context.Context) (Telemetry, error) { + if store == nil { + return Telemetry{}, errors.New("client spool is unavailable") + } + store.mu.Lock() + defer store.mu.Unlock() + if store.db == nil { + return Telemetry{}, errors.New("client spool is closed") + } + var result Telemetry + if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM commands WHERE terminal = 0 AND phase IN (?, ?)`, uint32(rvboxv1.CommandLifecycle_COMMAND_QUEUED), uint32(rvboxv1.CommandLifecycle_COMMAND_DISPATCHED)).Scan(&result.QueuedCommands); err != nil { + return Telemetry{}, err + } + if err := store.db.QueryRowContext(ctx, `SELECT client_total_charged_bytes FROM spool_counters WHERE singleton = 1`).Scan(&result.ChargedBytes); err != nil { + return Telemetry{}, err + } + return result, nil +} diff --git a/internal/client/spool/telemetry_test.go b/internal/client/spool/telemetry_test.go new file mode 100644 index 0000000..2f8a811 --- /dev/null +++ b/internal/client/spool/telemetry_test.go @@ -0,0 +1,22 @@ +package spool + +import ( + "context" + "path/filepath" + "testing" + "time" +) + +func TestTelemetryReportsQueuedDurableSpool_HP_OPS_07(t *testing.T) { + t.Parallel() + ctx := context.Background() + store := openTestStore(t, ctx, filepath.Join(t.TempDir(), "spool"), DefaultTombstoneLimit) + issue := testUUID(t, "019c46f1-1d02-7000-8000-0000000000d2") + if _, err := store.AcceptCommand(ctx, testCommand(issue, []byte("telemetry")), time.Now().UTC()); err != nil { + t.Fatal(err) + } + telemetry, err := store.Telemetry(ctx) + if err != nil || telemetry.QueuedCommands != 1 || telemetry.ChargedBytes == 0 { + t.Fatalf("telemetry = %#v, %v", telemetry, err) + } +} diff --git a/internal/server/store/telemetry.go b/internal/server/store/telemetry.go new file mode 100644 index 0000000..2031697 --- /dev/null +++ b/internal/server/store/telemetry.go @@ -0,0 +1,34 @@ +package store + +import ( + "context" + "errors" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +// Telemetry is the bounded operational state that the server publishes as +// gauges. It deliberately contains no client or command identifier. +type Telemetry struct { + QueuedCommands uint64 + ChargedBytes uint64 +} + +func (store *Store) Telemetry(ctx context.Context) (Telemetry, error) { + if store == nil { + return Telemetry{}, errors.New("server store is unavailable") + } + store.mu.Lock() + defer store.mu.Unlock() + if store.db == nil { + return Telemetry{}, ErrStoreClosed + } + var result Telemetry + if err := store.db.QueryRowContext(ctx, `SELECT count(*) FROM commands WHERE lifecycle = ?`, uint32(rvboxv1.CommandLifecycle_COMMAND_QUEUED)).Scan(&result.QueuedCommands); err != nil { + return Telemetry{}, err + } + if err := store.db.QueryRowContext(ctx, `SELECT command_charged_bytes FROM storage_counters WHERE singleton = 1`).Scan(&result.ChargedBytes); err != nil { + return Telemetry{}, err + } + return result, nil +} diff --git a/internal/server/store/telemetry_test.go b/internal/server/store/telemetry_test.go new file mode 100644 index 0000000..b2988c5 --- /dev/null +++ b/internal/server/store/telemetry_test.go @@ -0,0 +1,35 @@ +package store + +import ( + "context" + "crypto/sha256" + "path/filepath" + "testing" + "time" + + "github.com/rvbox/rvbox/internal/domain" +) + +func TestTelemetryReportsQueuedDurableState_HP_OPS_06(t *testing.T) { + t.Parallel() + ctx := context.Background() + persistence, err := Open(ctx, Options{DataDir: filepath.Join(t.TempDir(), "state"), BusyTimeout: time.Second}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = persistence.Close() }) + if _, err := persistence.RegisterClientSession(ctx, ClientRegistration{ClientID: "telemetry-client", Platform: 2, Architecture: "amd64", DaemonVersion: "test", DaemonCWD: `C:\`, SupportedShells: []byte{1}, ClientInstanceID: [16]byte{1}, SessionID: [16]byte{2}, ConnectedAt: time.Now().UTC()}); err != nil { + t.Fatal(err) + } + issue, err := domain.ParseUUIDv7("019c46f1-1d02-7000-8000-0000000000d1") + if err != nil { + t.Fatal(err) + } + if _, err := persistence.QueueCommand(ctx, QueueCommandInput{IssueUUID: issue, ClientID: "telemetry-client", IssueTime: time.Now().UTC(), ReceiptTime: time.Now().UTC(), ImmutableSHA256: sha256.Sum256([]byte("telemetry")), ExecutionSpec: []byte("spec")}); err != nil { + t.Fatal(err) + } + telemetry, err := persistence.Telemetry(ctx) + if err != nil || telemetry.QueuedCommands != 1 || telemetry.ChargedBytes == 0 { + t.Fatalf("telemetry = %#v, %v", telemetry, err) + } +} diff --git a/test/coverage.toml b/test/coverage.toml index 7a7d41e..94c1ecc 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -646,6 +646,15 @@ tests = [ "internal/server/session/agent_server_test.go:TestAgentServerEmitsBoundedRegistrationMetrics_HP_OPS_05", ] +[[requirements]] +id = "HP-OPS-06" +layer = "unit" +status = "implemented" +tests = [ + "internal/server/store/telemetry_test.go:TestTelemetryReportsQueuedDurableState_HP_OPS_06", + "internal/client/spool/telemetry_test.go:TestTelemetryReportsQueuedDurableSpool_HP_OPS_07", +] + [[requirements]] id = "HP-DISPATCH-10" layer = "unit"