feat: publish durable queue and spool gauges

This commit is contained in:
2026-09-11 09:28:54 +00:00
parent ea8a4d8846
commit 0e18af02bb
7 changed files with 178 additions and 0 deletions
+22
View File
@@ -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:
}
}
}
+22
View File
@@ -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),
+34
View File
@@ -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
}
+22
View File
@@ -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)
}
}
+34
View File
@@ -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
}
+35
View File
@@ -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)
}
}
+9
View File
@@ -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"