From 0ce3baa1774a100c4c9d5161c2a16dcd151ce116 Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Fri, 25 Sep 2026 04:24:39 +0800 Subject: [PATCH 1/7] Add bounded Core metrics store reads --- .../internal/db/queries/core_metrics.sql | 37 ++++ .../internal/db/sqlc/core_metrics.sql.go | 148 +++++++++++++++ .../agents-api/internal/store/core_metrics.go | 110 +++++++++++ .../internal/store/core_metrics_test.go | 176 ++++++++++++++++++ .../000068_core_metrics_indexes.sql | 9 + 5 files changed, 480 insertions(+) create mode 100644 services/agents-api/internal/db/queries/core_metrics.sql create mode 100644 services/agents-api/internal/db/sqlc/core_metrics.sql.go create mode 100644 services/agents-api/internal/store/core_metrics.go create mode 100644 services/agents-api/internal/store/core_metrics_test.go create mode 100644 services/agents-api/migrations/000068_core_metrics_indexes.sql diff --git a/services/agents-api/internal/db/queries/core_metrics.sql b/services/agents-api/internal/db/queries/core_metrics.sql new file mode 100644 index 000000000..19cde907a --- /dev/null +++ b/services/agents-api/internal/db/queries/core_metrics.sql @@ -0,0 +1,37 @@ +-- name: CoreExecutionSnapshot :one +SELECT count(*) FILTER (WHERE t.status = 'queued')::bigint AS queued_turns, + count(*) FILTER (WHERE t.status = 'queued' AND + (b.device_id IS NULL OR NOT (b.device_id = ANY(sqlc.arg(connected_device_ids)::uuid[]))))::bigint AS waiting_for_daemon, + count(*) FILTER (WHERE t.status = 'in_progress')::bigint AS in_progress_turns, + COALESCE(GREATEST(0, extract(epoch FROM (sqlc.arg(observed_at)::timestamptz - + min(t.created_at) FILTER (WHERE t.status = 'queued')))), 0)::double precision AS oldest_queued_seconds +FROM turns t LEFT JOIN session_devices b ON b.session_id = t.session_id +WHERE t.status IN ('queued', 'in_progress'); + +-- name: CoreInterruptedTurns :one +SELECT count(*)::bigint FROM turns +WHERE status = 'failed' AND outcome->>'error_code' = 'execution_interrupted' + AND completed_at >= sqlc.arg(range_start)::timestamptz + AND completed_at < sqlc.arg(range_end)::timestamptz; + +-- name: CoreQueueWaitSummary :one +SELECT count(*)::bigint AS samples, + COALESCE(percentile_cont(0.50) WITHIN GROUP (ORDER BY extract(epoch FROM (started_at - created_at)) * 1000), 0)::double precision AS p50_ms, + COALESCE(percentile_cont(0.95) WITHIN GROUP (ORDER BY extract(epoch FROM (started_at - created_at)) * 1000), 0)::double precision AS p95_ms +FROM turns +WHERE started_at >= sqlc.arg(range_start)::timestamptz + AND started_at < sqlc.arg(range_end)::timestamptz; + +-- name: CoreQueueWaitBuckets :many +SELECT floor(extract(epoch FROM (started_at - sqlc.arg(range_start)::timestamptz)) / + sqlc.arg(resolution_seconds)::integer)::integer AS bucket_number, + count(*)::bigint AS samples, + COALESCE(percentile_cont(0.95) WITHIN GROUP (ORDER BY extract(epoch FROM (started_at - created_at)) * 1000), 0)::double precision AS p95_ms +FROM turns +WHERE started_at >= sqlc.arg(range_start)::timestamptz + AND started_at < sqlc.arg(range_end)::timestamptz +GROUP BY bucket_number +ORDER BY bucket_number; + +-- name: CoreDatabaseSize :one +SELECT pg_database_size(current_database())::bigint; diff --git a/services/agents-api/internal/db/sqlc/core_metrics.sql.go b/services/agents-api/internal/db/sqlc/core_metrics.sql.go new file mode 100644 index 000000000..46fb6f12d --- /dev/null +++ b/services/agents-api/internal/db/sqlc/core_metrics.sql.go @@ -0,0 +1,148 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.29.0 +// source: core_metrics.sql + +package sqlc + +import ( + "context" + + "github.com/jackc/pgx/v5/pgtype" +) + +const coreDatabaseSize = `-- name: CoreDatabaseSize :one +SELECT pg_database_size(current_database())::bigint +` + +func (q *Queries) CoreDatabaseSize(ctx context.Context) (int64, error) { + row := q.db.QueryRow(ctx, coreDatabaseSize) + var column_1 int64 + err := row.Scan(&column_1) + return column_1, err +} + +const coreExecutionSnapshot = `-- name: CoreExecutionSnapshot :one +SELECT count(*) FILTER (WHERE t.status = 'queued')::bigint AS queued_turns, + count(*) FILTER (WHERE t.status = 'queued' AND + (b.device_id IS NULL OR NOT (b.device_id = ANY($1::uuid[]))))::bigint AS waiting_for_daemon, + count(*) FILTER (WHERE t.status = 'in_progress')::bigint AS in_progress_turns, + COALESCE(GREATEST(0, extract(epoch FROM ($2::timestamptz - + min(t.created_at) FILTER (WHERE t.status = 'queued')))), 0)::double precision AS oldest_queued_seconds +FROM turns t LEFT JOIN session_devices b ON b.session_id = t.session_id +WHERE t.status IN ('queued', 'in_progress') +` + +type CoreExecutionSnapshotParams struct { + ConnectedDeviceIds []pgtype.UUID `json:"connected_device_ids"` + ObservedAt pgtype.Timestamptz `json:"observed_at"` +} + +type CoreExecutionSnapshotRow struct { + QueuedTurns int64 `json:"queued_turns"` + WaitingForDaemon int64 `json:"waiting_for_daemon"` + InProgressTurns int64 `json:"in_progress_turns"` + OldestQueuedSeconds float64 `json:"oldest_queued_seconds"` +} + +func (q *Queries) CoreExecutionSnapshot(ctx context.Context, arg CoreExecutionSnapshotParams) (CoreExecutionSnapshotRow, error) { + row := q.db.QueryRow(ctx, coreExecutionSnapshot, arg.ConnectedDeviceIds, arg.ObservedAt) + var i CoreExecutionSnapshotRow + err := row.Scan( + &i.QueuedTurns, + &i.WaitingForDaemon, + &i.InProgressTurns, + &i.OldestQueuedSeconds, + ) + return i, err +} + +const coreInterruptedTurns = `-- name: CoreInterruptedTurns :one +SELECT count(*)::bigint FROM turns +WHERE status = 'failed' AND outcome->>'error_code' = 'execution_interrupted' + AND completed_at >= $1::timestamptz + AND completed_at < $2::timestamptz +` + +type CoreInterruptedTurnsParams struct { + RangeStart pgtype.Timestamptz `json:"range_start"` + RangeEnd pgtype.Timestamptz `json:"range_end"` +} + +func (q *Queries) CoreInterruptedTurns(ctx context.Context, arg CoreInterruptedTurnsParams) (int64, error) { + row := q.db.QueryRow(ctx, coreInterruptedTurns, arg.RangeStart, arg.RangeEnd) + var column_1 int64 + err := row.Scan(&column_1) + return column_1, err +} + +const coreQueueWaitBuckets = `-- name: CoreQueueWaitBuckets :many +SELECT floor(extract(epoch FROM (started_at - $1::timestamptz)) / + $2::integer)::integer AS bucket_number, + count(*)::bigint AS samples, + COALESCE(percentile_cont(0.95) WITHIN GROUP (ORDER BY extract(epoch FROM (started_at - created_at)) * 1000), 0)::double precision AS p95_ms +FROM turns +WHERE started_at >= $1::timestamptz + AND started_at < $3::timestamptz +GROUP BY bucket_number +ORDER BY bucket_number +` + +type CoreQueueWaitBucketsParams struct { + RangeStart pgtype.Timestamptz `json:"range_start"` + ResolutionSeconds int32 `json:"resolution_seconds"` + RangeEnd pgtype.Timestamptz `json:"range_end"` +} + +type CoreQueueWaitBucketsRow struct { + BucketNumber int32 `json:"bucket_number"` + Samples int64 `json:"samples"` + P95Ms float64 `json:"p95_ms"` +} + +func (q *Queries) CoreQueueWaitBuckets(ctx context.Context, arg CoreQueueWaitBucketsParams) ([]CoreQueueWaitBucketsRow, error) { + rows, err := q.db.Query(ctx, coreQueueWaitBuckets, arg.RangeStart, arg.ResolutionSeconds, arg.RangeEnd) + if err != nil { + return nil, err + } + defer rows.Close() + items := []CoreQueueWaitBucketsRow{} + for rows.Next() { + var i CoreQueueWaitBucketsRow + if err := rows.Scan(&i.BucketNumber, &i.Samples, &i.P95Ms); err != nil { + return nil, err + } + items = append(items, i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const coreQueueWaitSummary = `-- name: CoreQueueWaitSummary :one +SELECT count(*)::bigint AS samples, + COALESCE(percentile_cont(0.50) WITHIN GROUP (ORDER BY extract(epoch FROM (started_at - created_at)) * 1000), 0)::double precision AS p50_ms, + COALESCE(percentile_cont(0.95) WITHIN GROUP (ORDER BY extract(epoch FROM (started_at - created_at)) * 1000), 0)::double precision AS p95_ms +FROM turns +WHERE started_at >= $1::timestamptz + AND started_at < $2::timestamptz +` + +type CoreQueueWaitSummaryParams struct { + RangeStart pgtype.Timestamptz `json:"range_start"` + RangeEnd pgtype.Timestamptz `json:"range_end"` +} + +type CoreQueueWaitSummaryRow struct { + Samples int64 `json:"samples"` + P50Ms float64 `json:"p50_ms"` + P95Ms float64 `json:"p95_ms"` +} + +func (q *Queries) CoreQueueWaitSummary(ctx context.Context, arg CoreQueueWaitSummaryParams) (CoreQueueWaitSummaryRow, error) { + row := q.db.QueryRow(ctx, coreQueueWaitSummary, arg.RangeStart, arg.RangeEnd) + var i CoreQueueWaitSummaryRow + err := row.Scan(&i.Samples, &i.P50Ms, &i.P95Ms) + return i, err +} diff --git a/services/agents-api/internal/store/core_metrics.go b/services/agents-api/internal/store/core_metrics.go new file mode 100644 index 000000000..f7c3ee812 --- /dev/null +++ b/services/agents-api/internal/store/core_metrics.go @@ -0,0 +1,110 @@ +package store + +import ( + "context" + "time" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/db/sqlc" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgtype" +) + +type CoreExecutionSnapshot struct { + QueuedTurns, WaitingForDaemon, InProgressTurns int64 + OldestQueuedSeconds *float64 +} + +type CoreQueueWait struct{ P50, P95 *float64 } + +type CoreQueueWaitBucket struct { + Start time.Time + P95MS *float64 +} + +type CoreExecutionHistory struct { + Interrupted int64 + QueueWaitMS CoreQueueWait + Buckets []CoreQueueWaitBucket +} + +// ReadCoreExecutionSnapshot counts persisted root Turns across the deployment. +// WaitingForDaemon is the queued subset whose Session binding is absent from the +// supplied live registry IDs; callers must distinguish an unavailable registry +// from an observed empty one before using this method. Age is in seconds, and is +// nil when the queue is empty. No Session deletion filter hides operational state. +func (s *Store) ReadCoreExecutionSnapshot(ctx context.Context, now time.Time, connectedDeviceIDs []string) (CoreExecutionSnapshot, error) { + ids := make([]pgtype.UUID, 0, len(connectedDeviceIDs)) + for _, value := range connectedDeviceIDs { + id, err := parseID(value) + if err != nil { + return CoreExecutionSnapshot{}, err + } + ids = append(ids, id) + } + row, err := s.queries.CoreExecutionSnapshot(ctx, sqlc.CoreExecutionSnapshotParams{ + ConnectedDeviceIds: ids, ObservedAt: pgtype.Timestamptz{Time: now, Valid: true}, + }) + if err != nil { + return CoreExecutionSnapshot{}, err + } + result := CoreExecutionSnapshot{QueuedTurns: row.QueuedTurns, WaitingForDaemon: row.WaitingForDaemon, InProgressTurns: row.InProgressTurns} + if row.QueuedTurns > 0 { + result.OldestQueuedSeconds = &row.OldestQueuedSeconds + } + return result, nil +} + +// ReadCoreExecutionHistory reads a bounded, repeatable read-only snapshot of root +// Turn history. The interval is [start,end), with complete epoch-aligned buckets. +// Interruptions use failed Turns' completed_at and execution_interrupted error +// code. Queue waits use started_at-created_at in milliseconds, grouped by +// started_at, including retained history of deleted Sessions. Counts of an empty +// interval are zero; percentile values without observations are nil. Read errors +// invalidate the entire result and must not be presented as measured zeros. +func (s *Store) ReadCoreExecutionHistory(ctx context.Context, start, end time.Time, resolution time.Duration) (CoreExecutionHistory, error) { + span := end.Sub(start) + if resolution < time.Second || resolution%time.Second != 0 || span <= 0 || span > 7*24*time.Hour || + span%resolution != 0 || span/resolution > 1008 || start.Nanosecond() != 0 || end.Nanosecond() != 0 || + start.Unix()%int64(resolution/time.Second) != 0 || end.Unix()%int64(resolution/time.Second) != 0 { + return CoreExecutionHistory{}, ErrInvalidInput + } + result := CoreExecutionHistory{Buckets: make([]CoreQueueWaitBucket, int(span/resolution))} + for i := range result.Buckets { + result.Buckets[i].Start = start.Add(time.Duration(i) * resolution).UTC() + } + first, last := pgtype.Timestamptz{Time: start, Valid: true}, pgtype.Timestamptz{Time: end, Valid: true} + err := pgx.BeginTxFunc(ctx, s.pool, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}, func(tx pgx.Tx) error { + q := s.queries.WithTx(tx) + var err error + result.Interrupted, err = q.CoreInterruptedTurns(ctx, sqlc.CoreInterruptedTurnsParams{RangeStart: first, RangeEnd: last}) + if err != nil { + return err + } + wait, err := q.CoreQueueWaitSummary(ctx, sqlc.CoreQueueWaitSummaryParams{RangeStart: first, RangeEnd: last}) + if err != nil { + return err + } + if wait.Samples > 0 { + result.QueueWaitMS = CoreQueueWait{P50: &wait.P50Ms, P95: &wait.P95Ms} + } + rows, err := q.CoreQueueWaitBuckets(ctx, sqlc.CoreQueueWaitBucketsParams{RangeStart: first, RangeEnd: last, ResolutionSeconds: int32(resolution / time.Second)}) + if err != nil { + return err + } + for _, row := range rows { + if row.Samples > 0 { + result.Buckets[row.BucketNumber].P95MS = &row.P95Ms + } + } + return nil + }) + if err != nil { + return CoreExecutionHistory{}, err + } + return result, nil +} + +// ReadCoreDatabaseSize measures the current PostgreSQL database in bytes. +func (s *Store) ReadCoreDatabaseSize(ctx context.Context) (int64, error) { + return s.queries.CoreDatabaseSize(ctx) +} diff --git a/services/agents-api/internal/store/core_metrics_test.go b/services/agents-api/internal/store/core_metrics_test.go new file mode 100644 index 000000000..c208e8531 --- /dev/null +++ b/services/agents-api/internal/store/core_metrics_test.go @@ -0,0 +1,176 @@ +package store + +import ( + "context" + "errors" + "math" + "strings" + "testing" + "time" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" +) + +func coreMetricsSession(t *testing.T, pool *pgxpool.Pool, deleted bool) string { + t.Helper() + id := uuid.NewString() + _, err := pool.Exec(t.Context(), `INSERT INTO sessions(id,tenant_id,engine,idempotency_key,request_hash,deleted_at) + VALUES($1,$2,'codex','metrics','metrics',CASE WHEN $3::boolean THEN clock_timestamp() END)`, id, uuid.NewString(), deleted) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + for _, query := range []string{`DELETE FROM turns WHERE session_id=$1`, `DELETE FROM session_devices WHERE session_id=$1`, `DELETE FROM sessions WHERE id=$1`} { + if _, err := pool.Exec(context.Background(), query, id); err != nil { + t.Error(err) + } + } + }) + return id +} + +func coreMetricsTurn(t *testing.T, pool *pgxpool.Pool, session, status, code string, created time.Time, started, completed *time.Time) { + t.Helper() + _, err := pool.Exec(t.Context(), `INSERT INTO turns(id,session_id,status,created_at,started_at,completed_at,outcome) + VALUES($1,$2,$3,$4,$5,$6,jsonb_build_object('error_code',$7::text))`, uuid.NewString(), session, status, created, started, completed, code) + if err != nil { + t.Fatal(err) + } +} + +func TestCoreMetricsSnapshot(t *testing.T) { + s, pool := testStore(t) + now := time.Now().UTC() + baseline, err := s.ReadCoreExecutionSnapshot(t.Context(), now, []string{}) + if err != nil { + t.Fatal(err) + } + if baseline.QueuedTurns == 0 && baseline.OldestQueuedSeconds != nil { + t.Fatal("empty queue must not invent an age") + } + connected := uuid.NewString() + if _, err := pool.Exec(t.Context(), `INSERT INTO devices(id,tenant_id,name,credential_hash) VALUES($1,$2,'metrics',$3)`, connected, uuid.NewString(), strings.Repeat("a", 64)); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _, _ = pool.Exec(context.Background(), `DELETE FROM devices WHERE id=$1`, connected) }) + oldest := time.Date(1900, 1, 1, 0, 0, 0, 0, time.UTC) + for i, status := range []string{"queued", "queued", "queued", "in_progress", "waiting"} { + session := coreMetricsSession(t, pool, false) + coreMetricsTurn(t, pool, session, status, "", oldest.Add(time.Duration(i)*time.Second), nil, nil) + if i == 0 { + if _, err := pool.Exec(t.Context(), `INSERT INTO session_devices(session_id,device_id) VALUES($1,$2)`, session, connected); err != nil { + t.Fatal(err) + } + } + } + got, err := s.ReadCoreExecutionSnapshot(t.Context(), now, []string{connected}) + if err != nil { + t.Fatal(err) + } + if got.QueuedTurns != baseline.QueuedTurns+3 || got.InProgressTurns != baseline.InProgressTurns+1 || got.WaitingForDaemon != baseline.WaitingForDaemon+2 { + t.Fatalf("unexpected snapshot: %+v; baseline %+v", got, baseline) + } + if got.OldestQueuedSeconds == nil || math.Abs(*got.OldestQueuedSeconds-now.Sub(oldest).Seconds()) > 0.001 { + t.Fatalf("oldest age: %v", got.OldestQueuedSeconds) + } + disconnected, err := s.ReadCoreExecutionSnapshot(t.Context(), now, []string{}) + if err != nil || disconnected.WaitingForDaemon != got.WaitingForDaemon+1 { + t.Fatalf("registry disconnect: %+v, %v", disconnected, err) + } + if _, err := s.ReadCoreExecutionSnapshot(t.Context(), now, []string{"invalid-device-id"}); err == nil { + t.Fatal("invalid registry ID accepted") + } +} + +func TestCoreMetricsHistory(t *testing.T) { + s, pool := testStore(t) + start := time.Date(2001, 1, 1, 0, 0, 0, 0, time.UTC) + end := start.Add(3 * time.Minute) + // Boundary samples exercise inclusive start, exclusive end and exact buckets. + for i, offset := range []time.Duration{-time.Second, 0, 30 * time.Second, time.Minute, 3 * time.Minute} { + started := start.Add(offset) + wait := []time.Duration{9 * time.Second, 0, time.Second, 3 * time.Second, 9 * time.Second}[i] + completed := started.Add(time.Second) + if i == 4 { + completed = end + } + coreMetricsTurn(t, pool, coreMetricsSession(t, pool, i == 2), "failed", "execution_interrupted", started.Add(-wait), &started, &completed) + } + // Nonmatching outcomes, nonfailed status and a queued never-started Turn do + // not contribute an interruption; cancellation without a start is no wait. + for _, status := range []string{"completed", "cancelled", "failed"} { + completed := start.Add(20 * time.Second) + code := "execution_interrupted" + if status == "failed" { + code = "other_error" + } + coreMetricsTurn(t, pool, coreMetricsSession(t, pool, false), status, code, start, nil, &completed) + } + got, err := s.ReadCoreExecutionHistory(t.Context(), start, end, time.Minute) + if err != nil { + t.Fatal(err) + } + if got.Interrupted != 4 || len(got.Buckets) != 3 { + t.Fatalf("history counts: %+v", got) + } + check := func(label string, got *float64, want float64) { + t.Helper() + if got == nil || math.Abs(*got-want) > 0.001 { + t.Fatalf("%s: got %v, want %v", label, got, want) + } + } + check("range p50", got.QueueWaitMS.P50, 1000) + check("range p95", got.QueueWaitMS.P95, 2800) + check("first bucket p95", got.Buckets[0].P95MS, 950) + check("second bucket p95", got.Buckets[1].P95MS, 3000) + if got.Buckets[2].P95MS != nil || !got.Buckets[2].Start.Equal(start.Add(2*time.Minute)) { + t.Fatal("missing bucket must have aligned start and unknown percentile") + } + empty, err := s.ReadCoreExecutionHistory(t.Context(), end.Add(time.Hour), end.Add(2*time.Hour), time.Minute) + if err != nil || empty.Interrupted != 0 || empty.QueueWaitMS.P50 != nil || empty.QueueWaitMS.P95 != nil || len(empty.Buckets) != 60 { + t.Fatalf("empty history: %+v, %v", empty, err) + } + for _, bucket := range empty.Buckets { + if bucket.P95MS != nil { + t.Fatal("empty bucket invented zero percentile") + } + } + size, err := s.ReadCoreDatabaseSize(t.Context()) + if err != nil || size <= 0 { + t.Fatalf("database size: %d, %v", size, err) + } + ctx, cancel := context.WithCancel(t.Context()) + cancel() + if _, err := s.ReadCoreExecutionSnapshot(ctx, start, nil); err == nil { + t.Fatal("snapshot read error hidden") + } + if result, err := s.ReadCoreExecutionHistory(ctx, start, end, time.Minute); err == nil || result.Buckets != nil { + t.Fatal("history read error hidden or partial data returned") + } + if _, err := s.ReadCoreDatabaseSize(ctx); err == nil { + t.Fatal("database size read error hidden") + } +} + +func TestCoreMetricsHistoryBounds(t *testing.T) { + s := &Store{} + start := time.Date(2001, 1, 1, 0, 0, 0, 0, time.UTC) + for _, tc := range []struct { + start, end time.Time + step time.Duration + }{ + {start, start, time.Minute}, + {start, start.Add(time.Hour), 0}, + {start, start.Add(time.Hour), 1500 * time.Millisecond}, + {start, start.Add(8 * 24 * time.Hour), time.Hour}, + {start, start.Add(time.Hour), time.Second}, + {start.Add(time.Second), start.Add(time.Hour + time.Second), time.Minute}, + {start, start.Add(time.Hour + time.Second), time.Minute}, + {start.Add(time.Nanosecond), start.Add(time.Hour + time.Nanosecond), time.Minute}, + } { + if _, err := s.ReadCoreExecutionHistory(t.Context(), tc.start, tc.end, tc.step); !errors.Is(err, ErrInvalidInput) { + t.Fatalf("unbounded or unaligned range accepted: %+v, %v", tc, err) + } + } +} diff --git a/services/agents-api/migrations/000068_core_metrics_indexes.sql b/services/agents-api/migrations/000068_core_metrics_indexes.sql new file mode 100644 index 000000000..a45fa0f6a --- /dev/null +++ b/services/agents-api/migrations/000068_core_metrics_indexes.sql @@ -0,0 +1,9 @@ +-- +goose Up +CREATE INDEX turns_metrics_started_idx ON turns(started_at) INCLUDE (created_at) + WHERE started_at IS NOT NULL; +CREATE INDEX turns_metrics_interrupted_idx ON turns(completed_at) + WHERE status = 'failed' AND outcome->>'error_code' = 'execution_interrupted'; + +-- +goose Down +DROP INDEX turns_metrics_interrupted_idx; +DROP INDEX turns_metrics_started_idx; From d0b9103964340e8a8fdf6b8a5dea7247821b26a9 Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Fri, 25 Sep 2026 04:36:46 +0800 Subject: [PATCH 2/7] Observe worker slots ownership and scheduler polls --- .../agents-api/internal/execution/worker.go | 43 +++++-- .../internal/execution/worker_metrics.go | 111 ++++++++++++++++++ .../internal/execution/worker_metrics_test.go | 102 ++++++++++++++++ .../internal/execution/worker_schedule.go | 2 +- 4 files changed, 247 insertions(+), 11 deletions(-) create mode 100644 services/agents-api/internal/execution/worker_metrics.go create mode 100644 services/agents-api/internal/execution/worker_metrics_test.go diff --git a/services/agents-api/internal/execution/worker.go b/services/agents-api/internal/execution/worker.go index b182417a2..bd6339daa 100644 --- a/services/agents-api/internal/execution/worker.go +++ b/services/agents-api/internal/execution/worker.go @@ -11,8 +11,11 @@ import ( "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" ) +const workerSlotLimit = 4 + // Worker owns queued work; the database lease excludes a second execution service. type Worker struct { + metrics workerMetricsState dispatcher *Dispatcher admission *store.Store lease *store.ExecutionLease @@ -73,11 +76,16 @@ func StartWorker(ctx context.Context, dispatcher *Dispatcher) (*Worker, error) { _ = lease.Close(context.Background()) return nil, err } + worker.observeOwnership(nil) return worker, nil } // CheckOwnership checks the same database lease used for execution writes. -func (w *Worker) CheckOwnership(ctx context.Context) error { return w.lease.Ping(ctx) } +func (w *Worker) CheckOwnership(ctx context.Context) error { + err := w.lease.Ping(ctx) + w.observeOwnership(err) + return err +} func (w *Worker) SubmitInputs(ctx context.Context, tenant, session, key string, inputs []store.Input) ([]store.InputReceipt, error) { value, err := w.admission.GetSession(ctx, tenant, session) @@ -113,11 +121,12 @@ func (w *Worker) CreateSessionStream(ctx context.Context, tenant string, input s } // Run retains queued work across restarts, but never replays an uncertain claim. -func (w *Worker) Run(ctx context.Context) error { +func (w *Worker) Run(ctx context.Context) (runErr error) { defer w.stopOnce.Do(func() { close(w.stopped) }) ctx, cancel := context.WithCancel(ctx) var running sync.WaitGroup defer func() { + w.observeWorkerStop(runErr, ctx.Err()) cancel() if w.runtimes != nil { w.runtimes.stop() @@ -129,14 +138,15 @@ func (w *Worker) Run(ctx context.Context) error { } closeCtx, stop := context.WithTimeout(context.Background(), 5*time.Second) defer stop() - _ = w.lease.Close(closeCtx) + w.observeWorkerClosed(w.lease.Close(closeCtx)) }() active := make(map[string]bool) + w.observeSlots(len(active)) type completion struct { id string err error } - completed := make(chan completion, 4) + completed := make(chan completion, workerSlotLimit) lifecycleDone := make(chan error, 1) if w.runtimes != nil { running.Add(1) @@ -154,8 +164,8 @@ func (w *Worker) Run(ctx context.Context) error { request fileWriteRequest result fileWriteResult } - writesCompleted := make(chan writeCompletion, 4) - readsCompleted := make(chan readCompletion, 4) + writesCompleted := make(chan writeCompletion, workerSlotLimit) + readsCompleted := make(chan readCompletion, workerSlotLimit) reads := 0 ticker := time.NewTicker(250 * time.Millisecond) defer ticker.Stop() @@ -167,11 +177,12 @@ func (w *Worker) Run(ctx context.Context) error { case err := <-lifecycleDone: return err case request := <-w.fileWrites: - if request.ctx.Err() != nil || active[request.environment.SessionID] || len(active) == 4 { + if request.ctx.Err() != nil || active[request.environment.SessionID] || len(active) == workerSlotLimit { request.result <- fileWriteResult{err: ErrExecutionUnavailable} continue } active[request.environment.SessionID] = true + w.observeSlots(len(active)) running.Add(1) go func() { defer running.Done() @@ -179,15 +190,17 @@ func (w *Worker) Run(ctx context.Context) error { }() case write := <-writesCompleted: delete(active, write.request.environment.SessionID) + w.observeSlots(len(active)) write.request.result <- write.result case request := <-w.directoryReads: - if request.ctx.Err() != nil || reads == 4 || (!active[request.environment.SessionID] && len(active) == 4) { + if request.ctx.Err() != nil || reads == workerSlotLimit || (!active[request.environment.SessionID] && len(active) == workerSlotLimit) { request.reply(directoryReadResult{err: ErrExecutionUnavailable}) continue } reserved := !active[request.environment.SessionID] if reserved { active[request.environment.SessionID] = true + w.observeSlots(len(active)) } reads++ running.Add(1) @@ -204,37 +217,47 @@ func (w *Worker) Run(ctx context.Context) error { reads-- if read.id != "" { delete(active, read.id) + w.observeSlots(len(active)) } read.request.reply(read.result) case result := <-completed: delete(active, result.id) + w.observeSlots(len(active)) if result.err != nil { return result.err } case <-ticker.C: check, stop := context.WithTimeout(ctx, 5*time.Second) - err := w.lease.Ping(check) + err := w.CheckOwnership(check) stop() if err != nil { + w.observeSchedulerPoll(0, err) return err } if _, err := w.dispatcher.Store.ExpireEnvironmentInputs(ctx); err != nil { + w.observeSchedulerPoll(0, err) return err } if err := w.observeEnrolledRuntimes(ctx); err != nil { + w.observeSchedulerPoll(0, err) return err } - if len(active) == 4 { + if len(active) == workerSlotLimit { + w.observeSchedulerPoll(0, nil) continue } devices := w.dispatcher.Registry.Devices() if len(devices) == 0 { + w.observeSchedulerPoll(0, nil) continue } work, err := schedule.selectWork(ctx, w, devices, active) + w.observeSlots(len(active)) if err != nil { + w.observeSchedulerPoll(0, err) return err } + w.observeSchedulerPoll(len(work), nil) for _, item := range work { running.Add(1) go func() { diff --git a/services/agents-api/internal/execution/worker_metrics.go b/services/agents-api/internal/execution/worker_metrics.go new file mode 100644 index 000000000..d6bffc7fb --- /dev/null +++ b/services/agents-api/internal/execution/worker_metrics.go @@ -0,0 +1,111 @@ +package execution + +import ( + "errors" + "sync" + "time" +) + +// WorkerMetrics contains only observations from this worker's existing operations. +// Missing pointers mean the corresponding operation has not established a value. +type WorkerMetrics struct { + SlotsInUse *int64 + SlotsTotal *int64 + ExecutionOwner *bool + Scheduler WorkerJobMetrics +} + +// WorkerJobMetrics describes the last completed scheduling poll. Failed counts +// failed polls, not failed Turns; Processed is unknown when a poll fails. +type WorkerJobMetrics struct { + Status string + LastRunAt *time.Time + Processed *int64 + Failed *int64 +} + +type workerMetricsState struct { + mu sync.Mutex + value WorkerMetrics + closed bool +} + +// MetricsSnapshot copies bounded in-process observations without querying storage. +func (w *Worker) MetricsSnapshot() WorkerMetrics { + w.metrics.mu.Lock() + defer w.metrics.mu.Unlock() + value := w.metrics.value + value.SlotsInUse = copyMetric(value.SlotsInUse) + value.SlotsTotal = copyMetric(value.SlotsTotal) + value.ExecutionOwner = copyMetric(value.ExecutionOwner) + value.Scheduler.LastRunAt = copyMetric(value.Scheduler.LastRunAt) + value.Scheduler.Processed = copyMetric(value.Scheduler.Processed) + value.Scheduler.Failed = copyMetric(value.Scheduler.Failed) + if value.Scheduler.Status == "" { + value.Scheduler.Status = "unknown" + } + return value +} + +func copyMetric[T any](source *T) *T { + if source == nil { + return nil + } + value := *source + return &value +} + +func (w *Worker) observeSlots(active int) { + used, total := int64(active), int64(workerSlotLimit) + w.metrics.mu.Lock() + defer w.metrics.mu.Unlock() + w.metrics.value.SlotsInUse = &used + w.metrics.value.SlotsTotal = &total +} + +func (w *Worker) observeOwnership(err error) { + w.metrics.mu.Lock() + defer w.metrics.mu.Unlock() + if w.metrics.closed { + return + } + w.metrics.value.ExecutionOwner = nil + if err == nil { + owned := true + w.metrics.value.ExecutionOwner = &owned + } +} + +func (w *Worker) observeSchedulerPoll(processed int, err error) { + now, handled, failed := time.Now().UTC(), int64(processed), int64(0) + job := WorkerJobMetrics{Status: "ok", LastRunAt: &now, Processed: &handled, Failed: &failed} + if err != nil { + failed = 1 + job.Status, job.Processed = "failing", nil + } + w.metrics.mu.Lock() + defer w.metrics.mu.Unlock() + w.metrics.value.Scheduler = job +} + +func (w *Worker) observeWorkerStop(runErr, contextErr error) { + w.metrics.mu.Lock() + defer w.metrics.mu.Unlock() + w.metrics.value.Scheduler.Status = "stopped" + if runErr != nil && (contextErr == nil || !errors.Is(runErr, contextErr)) { + w.metrics.value.Scheduler.Status = "failing" + } +} + +func (w *Worker) observeWorkerClosed(err error) { + w.metrics.mu.Lock() + defer w.metrics.mu.Unlock() + w.metrics.closed = true + w.metrics.value.ExecutionOwner = nil + if err == nil { + owned := false + w.metrics.value.ExecutionOwner = &owned + } + used := int64(0) + w.metrics.value.SlotsInUse = &used +} diff --git a/services/agents-api/internal/execution/worker_metrics_test.go b/services/agents-api/internal/execution/worker_metrics_test.go new file mode 100644 index 000000000..e877ac0d4 --- /dev/null +++ b/services/agents-api/internal/execution/worker_metrics_test.go @@ -0,0 +1,102 @@ +package execution + +import ( + "context" + "errors" + "fmt" + "sync" + "testing" + "time" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" +) + +func TestWorkerMetricsUnknownAndDetached(t *testing.T) { + worker := &Worker{} + initial := worker.MetricsSnapshot() + if initial.SlotsInUse != nil || initial.SlotsTotal != nil || initial.ExecutionOwner != nil || initial.Scheduler.Status != "unknown" || initial.Scheduler.LastRunAt != nil || initial.Scheduler.Processed != nil || initial.Scheduler.Failed != nil { + t.Fatalf("unobserved worker reported values: %+v", initial) + } + worker.observeSlots(3) + worker.observeOwnership(nil) + worker.observeSchedulerPoll(2, nil) + first := worker.MetricsSnapshot() + if *first.SlotsInUse != 3 || *first.SlotsTotal != 4 || !*first.ExecutionOwner || *first.Scheduler.Processed != 2 || *first.Scheduler.Failed != 0 || first.Scheduler.Status != "ok" { + t.Fatalf("unexpected observation: %+v", first) + } + observedAt := *first.Scheduler.LastRunAt + *first.SlotsInUse, *first.SlotsTotal, *first.ExecutionOwner = 100, 100, false + *first.Scheduler.Processed, *first.Scheduler.Failed = 100, 100 + *first.Scheduler.LastRunAt = time.Time{} + second := worker.MetricsSnapshot() + if *second.SlotsInUse != 3 || *second.SlotsTotal != 4 || !*second.ExecutionOwner || *second.Scheduler.Processed != 2 || *second.Scheduler.Failed != 0 || !second.Scheduler.LastRunAt.Equal(observedAt) { + t.Fatal("caller mutation altered the worker's observations") + } + worker.observeSlots(1) + worker.observeSchedulerPoll(0, nil) + if *second.SlotsInUse != 3 || *second.Scheduler.Processed != 2 { + t.Fatal("new observations altered a retained snapshot") + } +} + +func TestWorkerMetricsFailuresAndClosure(t *testing.T) { + worker := &Worker{lease: &store.ExecutionLease{}} + worker.observeOwnership(nil) + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := worker.CheckOwnership(ctx); !errors.Is(err, context.Canceled) { + t.Fatalf("ownership error changed: %v", err) + } + if worker.MetricsSnapshot().ExecutionOwner != nil { + t.Fatal("failed ownership check reported a known lease owner") + } + failure := errors.New("poll failed") + worker.observeSchedulerPoll(3, failure) + job := worker.MetricsSnapshot().Scheduler + if job.Status != "failing" || job.Processed != nil || job.Failed == nil || *job.Failed != 1 || job.LastRunAt == nil { + t.Fatalf("failed poll concealed unknown progress: %+v", job) + } + worker.observeWorkerStop(failure, context.Canceled) + if worker.MetricsSnapshot().Scheduler.Status != "failing" { + t.Fatal("concurrent cancellation concealed the worker failure") + } + worker.observeWorkerStop(fmt.Errorf("stop: %w", context.Canceled), context.Canceled) + if worker.MetricsSnapshot().Scheduler.Status != "stopped" { + t.Fatal("normal cancellation reported as a failure") + } + worker.observeWorkerClosed(nil) + worker.observeOwnership(nil) + closed := worker.MetricsSnapshot() + if closed.ExecutionOwner == nil || *closed.ExecutionOwner || *closed.SlotsInUse != 0 { + t.Fatal("late ownership observation revived a closed worker") + } + uncertain := &Worker{} + uncertain.observeWorkerClosed(failure) + if uncertain.MetricsSnapshot().ExecutionOwner != nil { + t.Fatal("failed lease close reported certain ownership") + } +} + +func TestWorkerMetricsConcurrentSnapshots(t *testing.T) { + worker := &Worker{} + var running sync.WaitGroup + for i := 0; i < 4; i++ { + running.Add(1) + go func() { + defer running.Done() + for j := 0; j < 100; j++ { + worker.observeSlots(j % 5) + worker.observeOwnership(nil) + worker.observeSchedulerPoll(j%5, nil) + value := worker.MetricsSnapshot() + if *value.SlotsInUse < 0 || *value.SlotsInUse > *value.SlotsTotal || *value.SlotsTotal != 4 { + t.Errorf("inconsistent slot snapshot: %+v", value) + } + *value.SlotsInUse = -1 + *value.ExecutionOwner = false + *value.Scheduler.Processed = -1 + } + }() + } + running.Wait() +} diff --git a/services/agents-api/internal/execution/worker_schedule.go b/services/agents-api/internal/execution/worker_schedule.go index 013cab51f..450d417d0 100644 --- a/services/agents-api/internal/execution/worker_schedule.go +++ b/services/agents-api/internal/execution/worker_schedule.go @@ -46,7 +46,7 @@ func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []st } } var selected []scheduledWork - for len(turns)+len(environments) > 0 && len(active) < 4 { + for len(turns)+len(environments) > 0 && len(active) < workerSlotLimit { var item scheduledWork if len(environments) > 0 && (s.environmentFirst || len(turns) == 0) { value := environments[0] From b44cf846aa355cdafc3c7d3b33690de35b38fe0b Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Fri, 25 Sep 2026 04:38:40 +0800 Subject: [PATCH 3/7] Expose administrator Core operational metrics from existing owners --- CONTRIBUTING.md | 15 ++ contracts/agents-api/core-metrics.md | 99 ++++++++++ .../agents-api/sandbox-manager.openapi.yaml | 172 +++++++++++++++++ docs/api/web-management.md | 10 + scripts/build-agents-api-release.sh | 2 +- scripts/build-agents-api.sh | 8 +- scripts/build-core-distribution.sh | 2 +- .../agents-api/cmd/server/core_metrics.go | 100 ++++++++++ services/agents-api/cmd/server/main.go | 22 ++- .../agents-api/cmd/server/runtime_history.go | 11 +- services/agents-api/cmd/server/write_audit.go | 9 +- .../agents-api/cmd/server/write_audit_test.go | 2 +- .../internal/api/admin_resources.go | 1 + .../agents-api/internal/api/core_metrics.go | 66 +++++++ .../internal/api/core_metrics_test.go | 97 ++++++++++ services/agents-api/internal/api/errors.go | 3 + services/agents-api/internal/api/handler.go | 3 +- services/agents-api/internal/api/routing.go | 14 +- .../agents-api/internal/coremetrics/series.go | 70 +++++++ .../internal/coremetrics/service.go | 177 ++++++++++++++++++ .../internal/coremetrics/service_test.go | 161 ++++++++++++++++ .../agents-api/internal/coremetrics/types.go | 105 +++++++++++ .../runtimehistory/postgresreader/reader.go | 10 +- .../postgresreader/reader_test.go | 4 +- .../store/runtime_history_acceptance_test.go | 4 +- 25 files changed, 1144 insertions(+), 23 deletions(-) create mode 100644 contracts/agents-api/core-metrics.md create mode 100644 services/agents-api/cmd/server/core_metrics.go create mode 100644 services/agents-api/internal/api/core_metrics.go create mode 100644 services/agents-api/internal/api/core_metrics_test.go create mode 100644 services/agents-api/internal/coremetrics/series.go create mode 100644 services/agents-api/internal/coremetrics/service.go create mode 100644 services/agents-api/internal/coremetrics/service_test.go create mode 100644 services/agents-api/internal/coremetrics/types.go diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 23bb57af9..91fcfabc5 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -195,6 +195,21 @@ check targets remain; the separate `packages/codex-harness`, its build/check scr and its CI/`make check` gate are retired. Historical remote native probes are not current validation entrypoints. +## Core operational metrics + +The administrator-only `/core/v1/admin/core-metrics` contract is documented in +[core-metrics.md](contracts/agents-api/core-metrics.md). Keep this separate from +Agent outcome and Sandbox capacity views. Instrument existing worker and job +owners without changing scheduling, lease or retention behavior. Periodic pool +pings and bounded in-process samples have explicit restart gaps; unknown values +must remain null. Complete UTC buckets exclude the active partial bucket. Root +Turn history is queried read-only from PostgreSQL with native timestamps. +Count `execution_unavailable` at the existing HTTP error writer, once per rejected +response; never record request/response bodies or infer this count from every +503 or failed Turn. Builds inject the source commit with ldflags. No new monitoring +service or storage system is required. Keep the frontend response shape aligned +with the paired console contract. + ## Architecture boundaries The following execution rules are retained from the source contributor guide. diff --git a/contracts/agents-api/core-metrics.md b/contracts/agents-api/core-metrics.md new file mode 100644 index 000000000..aa8d2f20b --- /dev/null +++ b/contracts/agents-api/core-metrics.md @@ -0,0 +1,99 @@ +# Core operational metrics + +`GET /core/v1/admin/core-metrics?range=1h|6h|24h|7d` is a deployment-administrator +read. It implements the response shape agreed with Core Web PR #96 +(`53dc9d646bc6e3cc2cd9b8bbb353bc53c3513ecf`). It changes neither public `/v1` +resources nor Agent or Sandbox metrics. Project API keys cannot call it. + +Only `range` is accepted, once; omission defaults to `1h`. Empty, repeated, +unsupported or other query parameters return `400 invalid_request`. A missing +metrics service returns `503 core_metrics_unavailable`. Partial measurement +failures return the usual `200 core.metrics` envelope with `service.status` set +to `degraded` and unavailable fields set to JSON null. No database or native +error text, credentials, bodies, resource IDs or tenant labels are exposed. + +## Time and missing data + +The response's `range` uses UTC RFC 3339 boundaries. Its exclusive `end` is the +most recent complete bucket boundary. The included interval is `[start,end)`; +the current partial bucket is excluded from range aggregates and series. + +| Range | Bucket size | Buckets | +| --- | --- | --- | +| 1h | 60 seconds | 60 | +| 6h | 300 seconds | 72 | +| 24h | 900 seconds | 96 | +| 7d | 7200 seconds | 84 | + +Current gauges and range aggregates are intentionally different: current worker, +connection pool and process values are read when requested; queue gauges, database +size and deployment maintenance are sampled every 30 seconds with bounded I/O. +Samples older than 60 seconds are not reported as current. Queue, running and +pool series report the highest **observed** value in each bucket, not a claim +that all intermediate peaks were captured. Missing observations and the process's +partial first bucket stay null. Successful periodic ping samples produce linear +interpolated p50/p95; there is no request-triggered ping. + +A fixed-size in-process ring retains up to seven days of 30-second samples and +rejection counts. Restart loses those measurements: no synthetic backfill occurs. +The `execution.unavailable` count is null if the requested interval starts before +this process's observation began; an entirely observed interval with no rejections +is zero. PostgreSQL Turn history remains queryable across process restarts. +An empty queue has a measured count of zero but no oldest age. No started Turns +or successful ping samples means null percentiles, not zero latency. + +## Sources + +The envelope contains `object`, `range`, `service`, `execution`, `database`, +`jobs` and `process`, matching the typed client from PR #96. All numeric values +and `service.execution_owner` are nullable; lists of complete buckets are always +present. + +- `service.revision` is a full source commit injected into `main.buildRevision` + by the standalone builder's `-ldflags`. Manual builds without a valid revision + report null. `started_at` records process initialization. `execution_owner` + reflects the execution worker's existing lease checks, with unknown ownership + represented as null. `maintenance` means the saved deployment is in maintenance; + measurement or job failures take precedence as `degraded`. +- `execution.slots_in_use` is the worker's active Session reservation set. Its + configured capacity is four; environment input, Turns and file work share it. + It does not count native harness subprocesses. Worker-disabled installations + have zero configured execution slots. +- `queued_turns` and `in_progress_turns` count root rows in `turns`, including + operational state retained for deleted Sessions. Native Subagent views and + pending Environment input reservations are not extra queued root Turns. + `waiting_for_daemon` is the queued subset whose Session device binding is + absent from the actual connected-device registry. `oldest_queued_seconds` + measures the oldest queued row's `created_at`. +- `queue_wait_ms` uses `started_at - created_at`, in milliseconds, for Turns + started in the interval. Each bucket uses its own started Turns, with PostgreSQL + `percentile_cont`. `interrupted` counts failed Turns whose outcome error code is + `execution_interrupted`, using `completed_at` in the interval. +- `unavailable` counts actual HTTP errors emitted with code + `execution_unavailable`, once per rejected response. Other 503 codes and errors + occurring after a stream has already started are not counted. The existing error + writer reports the code; no response/request body capture is involved. +- `database.ping_ms` measures a periodic pool ping, including connection acquisition. + Pool `in_use`, `idle` and `max` come from `pgxpool.Stat()`. `size_bytes` is + `pg_database_size(current_database())`, not host disk usage. Failure to measure + one value does not turn it into zero. +- `process.memory_bytes` is Go `runtime.MemStats.Alloc` (allocated heap bytes), + not RSS or container memory. `goroutines` is `runtime.NumGoroutine()`. + +## Background jobs + +The four bounded IDs are `scheduler`, `runtime_sampler`, `history_cleanup` and +`audit_cleanup`. Each reports `status`, `last_run_at`, `processed` and `failed`. +A not-yet-observed run is unknown; a disabled or stopped loop is stopped. The +last-run time is the completion/observation of the last pass, not its next deadline. + +Scheduler processed counts selected Turn/environment work in that poll. Runtime +sampling counts observed and failed targets from its existing sweep result. +Cleanup processed counts confirmed removed rows; an unsuccessful cleanup reports +unknown processed count. Cleanup/scheduler failure counts identify failed passes, +not guessed numbers of lost rows or failed Turns. The actual scheduling, sampling, +retention and execution lifecycles keep their existing owners and timing. + +No additional telemetry database, monitoring server, model-provider probe, +message queue, object store, host disk measurement or scheduling mechanism is +introduced. Metrics cannot authorize execution or change resource ownership. diff --git a/contracts/agents-api/sandbox-manager.openapi.yaml b/contracts/agents-api/sandbox-manager.openapi.yaml index 252618239..05baa293f 100644 --- a/contracts/agents-api/sandbox-manager.openapi.yaml +++ b/contracts/agents-api/sandbox-manager.openapi.yaml @@ -172,6 +172,141 @@ definitions: $ref: '#/definitions/store.RuntimeNode' type: array type: object + coremetrics.Database: + properties: + ping_ms: + $ref: '#/definitions/coremetrics.Latency' + pool: + $ref: '#/definitions/coremetrics.Pool' + series: + items: + $ref: '#/definitions/coremetrics.DatabaseBucket' + type: array + size_bytes: + type: integer + type: object + coremetrics.DatabaseBucket: + properties: + ping_p95_ms: + type: number + pool_in_use: + type: integer + start: + type: string + type: object + coremetrics.Execution: + properties: + connected_daemons: + type: integer + in_progress_turns: + type: integer + interrupted: + type: integer + oldest_queued_seconds: + type: number + queue_wait_ms: + $ref: '#/definitions/coremetrics.Latency' + queued_turns: + type: integer + series: + items: + $ref: '#/definitions/coremetrics.ExecutionBucket' + type: array + slots_in_use: + type: integer + slots_total: + type: integer + unavailable: + type: integer + waiting_for_daemon: + type: integer + type: object + coremetrics.ExecutionBucket: + properties: + in_progress: + type: integer + queue_wait_p95_ms: + type: number + queued: + type: integer + start: + type: string + type: object + coremetrics.Job: + properties: + failed: + type: integer + id: + type: string + last_run_at: + type: string + processed: + type: integer + status: + type: string + type: object + coremetrics.Latency: + properties: + p50: + type: number + p95: + type: number + type: object + coremetrics.Pool: + properties: + idle: + type: integer + in_use: + type: integer + max: + type: integer + type: object + coremetrics.Process: + properties: + goroutines: + type: integer + memory_bytes: + type: integer + type: object + coremetrics.Range: + properties: + end: + type: string + resolution_seconds: + type: integer + start: + type: string + type: object + coremetrics.ServiceState: + properties: + execution_owner: + type: boolean + revision: + type: string + started_at: + type: string + status: + type: string + type: object + coremetrics.View: + properties: + database: + $ref: '#/definitions/coremetrics.Database' + execution: + $ref: '#/definitions/coremetrics.Execution' + jobs: + items: + $ref: '#/definitions/coremetrics.Job' + type: array + object: + type: string + process: + $ref: '#/definitions/coremetrics.Process' + range: + $ref: '#/definitions/coremetrics.Range' + service: + $ref: '#/definitions/coremetrics.ServiceState' + type: object store.AdminAssetCounts: properties: agents: @@ -2636,6 +2771,43 @@ paths: summary: Copy assets into another Project tags: - Core Administration + /core/v1/admin/core-metrics: + get: + description: Deployment administrator only. Complete UTC buckets; unknown measurements are null. Samples are process-local and are not backfilled after a restart. + parameters: + - description: Time range (default 1h) + enum: + - 1h + - 6h + - 24h + - 7d + in: query + name: range + type: string + produces: + - application/json + responses: + "200": + description: OK + schema: + $ref: '#/definitions/coremetrics.View' + "400": + description: Bad Request + schema: + $ref: '#/definitions/v1.ErrorResponse' + "401": + description: Unauthorized + schema: + $ref: '#/definitions/v1.ErrorResponse' + "503": + description: Service Unavailable + schema: + $ref: '#/definitions/v1.ErrorResponse' + security: + - DeploymentAdminAuth: [] + summary: Retrieve Core operational metrics + tags: + - Core Administration /core/v1/admin/projects: get: parameters: diff --git a/docs/api/web-management.md b/docs/api/web-management.md index 9494ad18a..ea82837d8 100644 --- a/docs/api/web-management.md +++ b/docs/api/web-management.md @@ -88,3 +88,13 @@ Origin/Host requests are rejected; unavailable Core or rejected upstream redirec return 502. Console authentication uses its own error envelope. Management resource errors and deletion constraints are documented in the administrator reference and generated schema; do not interpret every empty or failed read as an absent resource. + +## Core metrics + +`GET /core/v1/admin/core-metrics?range=1h|6h|24h|7d` returns Core process, execution +queue/slots, PostgreSQL and background-job measurements. It uses deployment +administrator authentication, rejects arbitrary query filters and never grants +Agent execution access. See the [exact measurement contract](../../contracts/agents-api/core-metrics.md) +for complete buckets, null values, units and process-local retention. Frontend +implementation is maintained separately; this backend change does not modify +Agent metrics or the public Agent API. diff --git a/scripts/build-agents-api-release.sh b/scripts/build-agents-api-release.sh index ca9e04cf5..abe6e64e4 100755 --- a/scripts/build-agents-api-release.sh +++ b/scripts/build-agents-api-release.sh @@ -55,7 +55,7 @@ if [[ "$go_version" != "$required_go" ]]; then printf 'Agents API release requires %s; found %s\n' "$required_go" "$go_version" >&2 exit 1 fi -AGENTS_API_BUILD_DIR="$release_context/package/bin" \ +AGENTS_API_BUILD_REVISION="$source_revision" AGENTS_API_BUILD_DIR="$release_context/package/bin" \ "$release_context/source/scripts/build-agents-api.sh" require_clean_source if [[ "$(git -C "$repo_root" rev-parse HEAD)" != "$source_revision" ]]; then diff --git a/scripts/build-agents-api.sh b/scripts/build-agents-api.sh index ed2f49add..4c2800080 100755 --- a/scripts/build-agents-api.sh +++ b/scripts/build-agents-api.sh @@ -11,6 +11,12 @@ for directory in "$runtime_root" "$output_dir"; do fi done +revision="${AGENTS_API_BUILD_REVISION:-$(git -C "$repo_root" rev-parse HEAD 2>/dev/null || true)}" +if [[ -n "$revision" && ! "$revision" =~ ^[0-9a-f]{40}$ ]]; then + printf 'Invalid Agents API source revision\n' >&2 + exit 1 +fi + mkdir -p "$runtime_root/cache/agents-api-builds" build_context="$(mktemp -d "$runtime_root/cache/agents-api-builds/source.XXXXXX")" trap 'rm -rf "$build_context"' EXIT @@ -31,7 +37,7 @@ tar -C "$repo_root" -cf - \ artifact="agents-api-$command" if [[ "$command" == server ]]; then artifact=agents-api; fi if [[ "$command" == sandbox-node ]]; then artifact=parsar-sandbox-node; fi - go build -mod=readonly -trimpath -buildvcs=false \ + go build -mod=readonly -trimpath -buildvcs=false -ldflags "-X main.buildRevision=$revision" \ -o "$build_context/bin/$artifact" "./services/agents-api/cmd/$command" done ) diff --git a/scripts/build-core-distribution.sh b/scripts/build-core-distribution.sh index ad6ede28f..72655f63e 100755 --- a/scripts/build-core-distribution.sh +++ b/scripts/build-core-distribution.sh @@ -101,7 +101,7 @@ cp services/agents-api/deploy/codex/seccomp.json "$bundle/runtime/" cp LICENSE "$bundle/" cp -R site "$bundle/site" -AGENTS_API_BUILD_DIR="$stage/core/bin" scripts/build-agents-api.sh +AGENTS_API_BUILD_REVISION="$revision" AGENTS_API_BUILD_DIR="$stage/core/bin" scripts/build-agents-api.sh ( cd services/agents-api/tools/microsandbox-provider GOWORK=off CGO_ENABLED=1 go build -mod=readonly -trimpath \ diff --git a/services/agents-api/cmd/server/core_metrics.go b/services/agents-api/cmd/server/core_metrics.go new file mode 100644 index 000000000..cb6956ac5 --- /dev/null +++ b/services/agents-api/cmd/server/core_metrics.go @@ -0,0 +1,100 @@ +package main + +import ( + "context" + "time" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/coremetrics" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" + "github.com/jackc/pgx/v5/pgxpool" +) + +// Release builders set this full source commit with -ldflags. +var buildRevision string +var processStartedAt = time.Now().UTC() + +type coreMetricsSource struct { + store *store.Store + pool *pgxpool.Pool + worker *execution.Worker + registry *gateway.Registry +} + +func metricPtr[T any](value T) *T { return &value } +func (s *coreMetricsSource) Live() coremetrics.Live { + stat := s.pool.Stat() + live := coremetrics.Live{Pool: coremetrics.Pool{InUse: metricPtr(int64(stat.AcquiredConns())), Idle: metricPtr(int64(stat.IdleConns())), Max: metricPtr(int64(stat.MaxConns()))}, Scheduler: coremetrics.Job{ID: "scheduler", Status: "stopped"}, ExecutionOwner: metricPtr(false), SlotsTotal: metricPtr(int64(0)), SlotsInUse: metricPtr(int64(0))} + if s.registry != nil { + live.ConnectedDaemons = metricPtr(int64(len(s.registry.Devices()))) + } + if s.worker != nil { + w := s.worker.MetricsSnapshot() + live.SlotsInUse, live.SlotsTotal, live.ExecutionOwner = w.SlotsInUse, w.SlotsTotal, w.ExecutionOwner + live.Scheduler = coremetrics.Job{ID: "scheduler", Status: w.Scheduler.Status, LastRunAt: w.Scheduler.LastRunAt, Processed: w.Scheduler.Processed, Failed: w.Scheduler.Failed} + } + return live +} +func (s *coreMetricsSource) Sample(ctx context.Context) coremetrics.Sample { + sample := coremetrics.Sample{Healthy: true, PoolInUse: metricPtr(int64(s.pool.Stat().AcquiredConns()))} + start := time.Now() + pingCtx, cancel := context.WithTimeout(ctx, time.Second) + err := s.pool.Ping(pingCtx) + cancel() + if err == nil { + sample.PingMS = metricPtr(float64(time.Since(start)) / float64(time.Millisecond)) + } else { + sample.Healthy = false + } + devices := []string{} + if s.registry != nil { + devices = s.registry.Devices() + } + counts, err := s.store.ReadCoreExecutionSnapshot(ctx, time.Now(), devices) + if err != nil { + sample.Healthy = false + } else { + sample.Queued, sample.InProgress = metricPtr(counts.QueuedTurns), metricPtr(counts.InProgressTurns) + sample.OldestQueuedSeconds = counts.OldestQueuedSeconds + if s.registry != nil { + sample.WaitingForDaemon = metricPtr(counts.WaitingForDaemon) + } + } + size, err := s.store.ReadCoreDatabaseSize(ctx) + if err != nil { + sample.Healthy = false + } else { + sample.DatabaseSize = &size + } + deployment, err := s.store.GetRuntimeDeployment(ctx) + if err != nil { + sample.Healthy = false + } else { + sample.Maintenance = metricPtr(deployment.Maintenance) + } + return sample +} +func (s *coreMetricsSource) History(ctx context.Context, start, end time.Time, step time.Duration) (coremetrics.History, error) { + value, err := s.store.ReadCoreExecutionHistory(ctx, start, end, step) + if err != nil { + return coremetrics.History{}, err + } + result := coremetrics.History{Interrupted: value.Interrupted, QueueWaitMS: coremetrics.Latency{P50: value.QueueWaitMS.P50, P95: value.QueueWaitMS.P95}, Buckets: map[time.Time]*float64{}} + for _, bucket := range value.Buckets { + result.Buckets[bucket.Start.UTC()] = bucket.P95MS + } + return result, nil +} + +func reportCleanupResult(metrics *coremetrics.Service, job string, count int64, err error) { + if metrics == nil { + return + } + processed, failed := &count, int64(0) + if err != nil { + processed = nil + failed = 1 + } + metrics.ReportJob(job, time.Now(), processed, &failed, err) +} diff --git a/services/agents-api/cmd/server/main.go b/services/agents-api/cmd/server/main.go index db39a98a9..9d63a7647 100644 --- a/services/agents-api/cmd/server/main.go +++ b/services/agents-api/cmd/server/main.go @@ -33,6 +33,7 @@ import ( "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/gateway" "github.com/MiniMax-AI-Dev/parsar/internal/obs/log" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/api" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/coremetrics" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtime" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeenrollment" @@ -90,6 +91,9 @@ func run() error { return err } executionStore := store.NewWithCredentialCipherAndOAuthRefresh(pool, credentialKey, oauthClient) + metricsSource := &coreMetricsSource{store: executionStore, pool: pool} + metrics := coremetrics.New(processStartedAt, buildRevision, metricsSource) + metrics.StopJob("runtime_sampler") auth, err := api.NewDatabaseAuthenticator(executionStore) if err != nil { return err @@ -102,7 +106,7 @@ func run() error { auditCleanupDone := make(chan struct{}) go func() { defer close(auditCleanupDone) - runWriteAuditCleanup(auditCleanupCtx, executionStore, auditRetention) + runWriteAuditCleanup(auditCleanupCtx, executionStore, auditRetention, metrics) }() defer func() { cancelAuditCleanup(); <-auditCleanupDone }() var workerDone chan error @@ -158,10 +162,10 @@ func run() error { cleanupDone := make(chan struct{}) go func() { defer close(cleanupDone) - runHistoryCleanup(cleanupCtx, history.Prune) + runHistoryCleanup(cleanupCtx, history.Prune, metrics) }() defer func() { cancelCleanup(); <-cleanupDone }() - options := []api.Option{api.WithSubagents(executionStore), api.WithSkills(executionStore), api.WithSourceFiles(executionStore), api.WithSessionArtifacts(executionStore), api.WithRuntimeObservations(observationService)} + options := []api.Option{api.WithCoreMetrics(metrics), api.WithSubagents(executionStore), api.WithSkills(executionStore), api.WithSourceFiles(executionStore), api.WithSessionArtifacts(executionStore), api.WithRuntimeObservations(observationService)} if managedNodes != nil { options = append(options, api.WithSandboxManager(executionStore, managedNodes.admin)) } @@ -231,6 +235,11 @@ func run() error { sampler, err := runtimeobs.NewSampler(observationResolver, observationService, worker, runtimeobs.SamplerOptions{ Interval: history.SampleInterval, Report: func(result runtimeobs.SweepResult) { + var sampleErr error + if !result.Complete { + sampleErr = errors.New("incomplete Runtime sampling sweep") + } + metrics.ReportJob("runtime_sampler", result.CompletedAt, metricPtr(int64(result.Observed)), metricPtr(int64(result.Failed)), sampleErr) fields := []any{"listed", result.Listed, "observed", result.Observed, "failed", result.Failed, "complete", result.Complete} if result.Complete { log.Bg().Debug("Runtime history sampling sweep complete", fields...) @@ -244,12 +253,17 @@ func run() error { } samplerCtx, cancelSampler := context.WithCancel(ctx) samplerDone := make(chan error, 1) - go func() { samplerDone <- sampler.Run(samplerCtx) }() + go func() { defer metrics.StopJob("runtime_sampler"); samplerDone <- sampler.Run(samplerCtx) }() defer func() { cancelSampler() <-samplerDone }() } + metricsSource.worker, metricsSource.registry = worker, registry + metricsCtx, cancelMetrics := context.WithCancel(ctx) + metricsDone := make(chan struct{}) + go func() { defer close(metricsDone); metrics.Run(metricsCtx) }() + defer func() { cancelMetrics(); <-metricsDone }() startupManaged := managed if managedNodes != nil && managedNodes.setup != nil { startupManaged = managedNodes.setup.selected.Load() diff --git a/services/agents-api/cmd/server/runtime_history.go b/services/agents-api/cmd/server/runtime_history.go index 131c5261e..f08a7a650 100644 --- a/services/agents-api/cmd/server/runtime_history.go +++ b/services/agents-api/cmd/server/runtime_history.go @@ -11,6 +11,7 @@ import ( "strings" "time" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/coremetrics" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimehistory" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimehistory/postgresreader" "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/runtimeobs" @@ -43,7 +44,7 @@ type runtimeHistorySetup struct { Exporter runtimeHistoryExporter Reader runtimehistory.Reader SampleInterval time.Duration - Prune func(context.Context) error + Prune func(context.Context) (int64, error) } type runtimeHistoryExporter interface { @@ -121,13 +122,17 @@ func closeRuntimeHistory(ctx context.Context, exporter runtimeHistoryExporter) { } // Retention also runs without active Runtimes. Each bounded pass has its own deadline. -func runHistoryCleanup(ctx context.Context, prune func(context.Context) error) { +func runHistoryCleanup(ctx context.Context, prune func(context.Context) (int64, error), metrics *coremetrics.Service) { + if metrics != nil { + defer metrics.StopJob("history_cleanup") + } ticker := time.NewTicker(time.Minute) defer ticker.Stop() for { pruneCtx, cancel := context.WithTimeout(ctx, 2*time.Second) - _ = prune(pruneCtx) + count, err := prune(pruneCtx) cancel() + reportCleanupResult(metrics, "history_cleanup", count, err) select { case <-ctx.Done(): return diff --git a/services/agents-api/cmd/server/write_audit.go b/services/agents-api/cmd/server/write_audit.go index c1524cdc2..ebe20904c 100644 --- a/services/agents-api/cmd/server/write_audit.go +++ b/services/agents-api/cmd/server/write_audit.go @@ -7,6 +7,7 @@ import ( "time" "github.com/MiniMax-AI-Dev/parsar/internal/obs/log" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/coremetrics" ) type writeAuditPruner interface { @@ -25,13 +26,17 @@ func writeAuditRetention() (time.Duration, error) { return duration, nil } -func runWriteAuditCleanup(ctx context.Context, s writeAuditPruner, retention time.Duration) { +func runWriteAuditCleanup(ctx context.Context, s writeAuditPruner, retention time.Duration, metrics *coremetrics.Service) { + if metrics != nil { + defer metrics.StopJob("audit_cleanup") + } ticker := time.NewTicker(time.Minute) defer ticker.Stop() for { pruneCtx, cancel := context.WithTimeout(ctx, 5*time.Second) - _, err := s.DeleteExpiredWriteOperations(pruneCtx, time.Now().Add(-retention), 1000) + count, err := s.DeleteExpiredWriteOperations(pruneCtx, time.Now().Add(-retention), 1000) cancel() + reportCleanupResult(metrics, "audit_cleanup", count, err) if err != nil && ctx.Err() == nil { log.Ctx(ctx).Warn("Write audit retention cleanup failed") } diff --git a/services/agents-api/cmd/server/write_audit_test.go b/services/agents-api/cmd/server/write_audit_test.go index aaa3cad4e..e4863e3eb 100644 --- a/services/agents-api/cmd/server/write_audit_test.go +++ b/services/agents-api/cmd/server/write_audit_test.go @@ -44,7 +44,7 @@ func TestWriteAuditCleanupBoundedAndCancellable(t *testing.T) { defer cancel() probe := &auditPruneProbe{cancel: cancel} before := time.Now().Add(-24 * time.Hour) - runWriteAuditCleanup(ctx, probe, 24*time.Hour) + runWriteAuditCleanup(ctx, probe, 24*time.Hour, nil) if !probe.called || probe.limit != 1000 || !probe.deadline || probe.cutoff.Before(before) || probe.cutoff.After(time.Now().Add(-24*time.Hour)) { t.Fatalf("bad cleanup %+v", probe) } diff --git a/services/agents-api/internal/api/admin_resources.go b/services/agents-api/internal/api/admin_resources.go index 1e6dd59a0..b5724fde4 100644 --- a/services/agents-api/internal/api/admin_resources.go +++ b/services/agents-api/internal/api/admin_resources.go @@ -39,6 +39,7 @@ func (h *Handler) registerAdminResourceRoutes(router chi.Router) { } router.Group(func(r chi.Router) { + r.Get("/core-metrics", h.getCoreMetrics) r.Get("/startup-configuration", h.adminStartupConfiguration) r.Get("/runtime-history/capabilities", h.adminRuntimeHistoryCapabilities) if h.adminManagement != nil { diff --git a/services/agents-api/internal/api/core_metrics.go b/services/agents-api/internal/api/core_metrics.go new file mode 100644 index 000000000..414204d9d --- /dev/null +++ b/services/agents-api/internal/api/core_metrics.go @@ -0,0 +1,66 @@ +package api + +import ( + "context" + "net/http" + "net/url" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/coremetrics" +) + +type CoreMetricsService interface { + Read(context.Context, string) (coremetrics.View, error) + RecordUnavailable() +} + +func WithCoreMetrics(service CoreMetricsService) Option { + return func(h *Handler) { h.coreMetrics = service } +} + +// @Summary Retrieve Core operational metrics +// @Description Deployment administrator only. Complete UTC buckets; unknown measurements are null. Samples are process-local and are not backfilled after a restart. +// @Tags Core Administration +// @Produce json +// @Security DeploymentAdminAuth +// @Param range query string false "Time range (default 1h)" Enums(1h,6h,24h,7d) +// @Success 200 {object} coremetrics.View +// @Failure 400,401,503 {object} v1.ErrorResponse +// @Router /core/v1/admin/core-metrics [get] +func (h *Handler) getCoreMetrics(w http.ResponseWriter, r *http.Request) { + values, err := url.ParseQuery(r.URL.RawQuery) + if err != nil || len(values) > 1 || (len(values) == 1 && len(values["range"]) != 1) { + writeError(w, http.StatusBadRequest, "invalid_request", "Only one supported range parameter is allowed.") + return + } + name := "1h" + if v, ok := values["range"]; ok { + name = v[0] + } + switch name { + case "1h", "6h", "24h", "7d": + default: + writeError(w, http.StatusBadRequest, "invalid_request", "range must be 1h, 6h, 24h or 7d.") + return + } + if h.coreMetrics == nil { + writeError(w, http.StatusServiceUnavailable, "core_metrics_unavailable", "Core metrics are not configured.") + return + } + value, err := h.coreMetrics.Read(r.Context(), name) + if err != nil { + writeError(w, http.StatusServiceUnavailable, "core_metrics_unavailable", "Core metrics could not be read.") + return + } + writeJSON(w, http.StatusOK, value) +} + +func (h *Handler) responseHeaders(next http.Handler) http.Handler { + if h.coreMetrics == nil { + return agentsResponseHeaders(next) + } + return responseHeadersWithErrors(next, func(code string) { + if code == "execution_unavailable" { + h.coreMetrics.RecordUnavailable() + } + }) +} diff --git a/services/agents-api/internal/api/core_metrics_test.go b/services/agents-api/internal/api/core_metrics_test.go new file mode 100644 index 000000000..bff441d14 --- /dev/null +++ b/services/agents-api/internal/api/core_metrics_test.go @@ -0,0 +1,97 @@ +package api + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "testing" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/coremetrics" +) + +type metricsFixture struct { + calls, refusals int + name string + err error +} + +func (f *metricsFixture) Read(_ context.Context, name string) (coremetrics.View, error) { + f.calls++ + f.name = name + return coremetrics.View{Object: "core.metrics"}, f.err +} +func (f *metricsFixture) RecordUnavailable() { f.refusals++ } +func TestCoreMetricsAdministratorContract(t *testing.T) { + key := callerBinding() + auth, _ := NewAuthenticator([]APIKey{key}) + admin, _ := NewDeploymentAuthenticator([]string{device.HashCredential("admin")}) + f := &metricsFixture{} + h, err := NewHandler(&recordingStore{}, auth, "codex", WithProjectAPIKeys(managementProjectStore(key), admin), WithCoreMetrics(f)) + if err != nil { + t.Fatal(err) + } + path := "/core/v1/admin/core-metrics" + for _, token := range []string{"", "unknown", "caller"} { + w := projectKeyHTTP(h, "GET", path, token, "") + if w.Code != 401 || f.calls != 0 { + t.Fatal(w.Code, f.calls) + } + } + for _, name := range []string{"1h", "6h", "24h", "7d"} { + w := projectKeyHTTP(h, "GET", path+"?range="+name, "admin", "") + if w.Code != 200 || f.name != name { + t.Fatal(w.Code, f.name) + } + } + w := projectKeyHTTP(h, "GET", path, "admin", "") + if w.Code != 200 || f.name != "1h" { + t.Fatal(w.Code, f.name) + } + before := f.calls + for _, query := range []string{"?range=", "?range=1h&range=6h", "?project_id=x", "?range=1h&tenant_id=x", "?range=1h;secret=x", "?range=%xx", "?range=8d"} { + w := projectKeyHTTP(h, "GET", path+query, "admin", "") + if w.Code != 400 || f.calls != before { + t.Fatal(query, w.Code, f.calls) + } + } + if w := projectKeyHTTP(h, "POST", path, "admin", "{}"); w.Code != 405 { + t.Fatal(w.Code) + } + f.err = errors.New("private credentials") + if w := projectKeyHTTP(h, "GET", path, "admin", ""); w.Code != 503 || w.Body.String() == f.err.Error() { + t.Fatal(w.Code, w.Body.String()) + } +} +func TestCoreRejectionObservationPreservesResponsesAndFlush(t *testing.T) { + for _, code := range []string{"execution_unavailable", "environment_unavailable", "invalid_api_key"} { + f := &metricsFixture{} + h := &Handler{coreMetrics: f} + out := httptest.NewRecorder() + handler := h.responseHeaders(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + writeError(w, 503, code, "not recorded") + writeError(w, 503, code, "second write") + })) + handler.ServeHTTP(out, httptest.NewRequest("GET", "/v1/agents/sessions", nil)) + want := 0 + if code == "execution_unavailable" { + want = 1 + } + if f.refusals != want { + t.Fatal(code, f.refusals) + } + } + f := &metricsFixture{} + h := &Handler{coreMetrics: f} + out := httptest.NewRecorder() + h.responseHeaders(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if err := http.NewResponseController(w).Flush(); err != nil { + t.Fatal(err) + } + writeError(w, 503, "execution_unavailable", "already streaming") + })).ServeHTTP(out, httptest.NewRequest("GET", "/", nil)) + if !out.Flushed || f.refusals != 0 { + t.Fatal("stream response changed or counted as refusal") + } +} diff --git a/services/agents-api/internal/api/errors.go b/services/agents-api/internal/api/errors.go index 261ae0238..da60e8d30 100644 --- a/services/agents-api/internal/api/errors.go +++ b/services/agents-api/internal/api/errors.go @@ -24,6 +24,9 @@ func writeJSON(w http.ResponseWriter, status int, value any) { // Every 401 has type invalid_request_error, as every observed official 401 // does (HP-07); an empty code serializes as null. func writeError(w http.ResponseWriter, status int, code, message string, param ...string) { + if observer, ok := w.(interface{ reportAPIError(string) }); ok { + observer.reportAPIError(code) + } kind := "invalid_request_error" if status >= 500 { kind = "server_error" diff --git a/services/agents-api/internal/api/handler.go b/services/agents-api/internal/api/handler.go index 8f1c855f9..fc65c568c 100644 --- a/services/agents-api/internal/api/handler.go +++ b/services/agents-api/internal/api/handler.go @@ -37,6 +37,7 @@ type ResourceStore interface { } type Handler struct { + coreMetrics CoreMetricsService sandboxStore *store.Store deploymentAuth *DeploymentAuthenticator sandboxSetup func(context.Context, store.SandboxDeploymentSetupRequest) (store.RuntimeDeploymentView, error) @@ -85,7 +86,7 @@ func NewHandler(s ResourceStore, auth *Authenticator, engine string, options ... // Beta group, has the JSON body and Allow header. func (h *Handler) routes() *chi.Mux { router := chi.NewRouter() - router.Use(agentsResponseHeaders, log.HTTPMiddleware, middleware.GetHead) + router.Use(h.responseHeaders, log.HTTPMiddleware, middleware.GetHead) router.MethodNotAllowed(methodNotAllowed) router.Get("/healthz", func(w http.ResponseWriter, _ *http.Request) { writeJSON(w, http.StatusOK, map[string]string{"status": "ok"}) diff --git a/services/agents-api/internal/api/routing.go b/services/agents-api/internal/api/routing.go index fec1eebed..d89125d90 100644 --- a/services/agents-api/internal/api/routing.go +++ b/services/agents-api/internal/api/routing.go @@ -123,13 +123,17 @@ func cleanPath(p string) string { // Cache-Control extensions remain. Organization and project headers are not // reported: Core's project scope is configured, not account-derived. func agentsResponseHeaders(next http.Handler) http.Handler { + return responseHeadersWithErrors(next, nil) +} + +func responseHeadersWithErrors(next http.Handler, report func(string)) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { id := newRequestID() header := w.Header() header.Set("X-Request-Id", id) header.Set("Openai-Version", "2020-10-01") header.Set("X-Content-Type-Options", "nosniff") - writer := &processingTimeWriter{ResponseWriter: w, started: time.Now()} + writer := &processingTimeWriter{ResponseWriter: w, started: time.Now(), report: report} next.ServeHTTP(writer, r.WithContext(log.WithRequestID(r.Context(), id))) }) } @@ -147,6 +151,14 @@ type processingTimeWriter struct { http.ResponseWriter started time.Time stamped bool + report func(string) +} + +// reportAPIError observes the emitted code without reading or retaining bodies. +func (w *processingTimeWriter) reportAPIError(code string) { + if !w.stamped && w.report != nil { + w.report(code) + } } func (w *processingTimeWriter) stamp() { diff --git a/services/agents-api/internal/coremetrics/series.go b/services/agents-api/internal/coremetrics/series.go new file mode 100644 index 000000000..ccb524686 --- /dev/null +++ b/services/agents-api/internal/coremetrics/series.go @@ -0,0 +1,70 @@ +package coremetrics + +import ( + "math" + "sort" + "time" +) + +// Caller holds mu. Gauges are the highest observed value, not an interpolated +// history. Missing buckets, including the process's partial first bucket, stay null. +func (s *Service) series(view *View) { + step := time.Duration(view.Range.ResolutionSeconds) * time.Second + count := int(view.Range.End.Sub(view.Range.Start) / step) + view.Execution.Series = make([]ExecutionBucket, count) + view.Database.Series = make([]DatabaseBucket, count) + pings := make([][]float64, count) + var allPings []float64 + for i := 0; i < count; i++ { + start := view.Range.Start.Add(time.Duration(i) * step) + view.Execution.Series[i].Start = start + view.Database.Series[i].Start = start + } + for _, sample := range s.samples { + if sample.At.Before(view.Range.Start) || !sample.At.Before(view.Range.End) || sample.At.IsZero() { + continue + } + i := int(sample.At.Sub(view.Range.Start) / step) + if sample.PingMS != nil { + allPings = append(allPings, *sample.PingMS) + } + if view.Execution.Series[i].Start.Before(s.started) { + continue + } + maxInto(&view.Execution.Series[i].Queued, sample.Queued) + maxInto(&view.Execution.Series[i].InProgress, sample.InProgress) + maxInto(&view.Database.Series[i].PoolInUse, sample.PoolInUse) + if sample.PingMS != nil { + pings[i] = append(pings[i], *sample.PingMS) + } + } + for i := range pings { + view.Database.Series[i].PingP95MS = percentile(pings[i], .95) + } + view.Database.PingMS = Latency{P50: percentile(allPings, .5), P95: percentile(allPings, .95)} + // A lifetime counter cannot prove the missing part of a pre-start range. + if !view.Range.Start.Before(s.started) { + total := int64(0) + for _, slot := range s.refusals { + at := time.Unix(slot.tick*int64(SampleInterval/time.Second), 0) + if !at.Before(view.Range.Start) && at.Before(view.Range.End) { + total += slot.count + } + } + view.Execution.Unavailable = &total + } +} +func maxInto(target **int64, value *int64) { + if value != nil && (*target == nil || **target < *value) { + *target = ptr(*value) + } +} +func percentile(values []float64, quantile float64) *float64 { + if len(values) == 0 { + return nil + } + sort.Float64s(values) + n := float64(len(values)-1) * quantile + lo, hi := int(math.Floor(n)), int(math.Ceil(n)) + return ptr(values[lo] + (values[hi]-values[lo])*(n-float64(lo))) +} diff --git a/services/agents-api/internal/coremetrics/service.go b/services/agents-api/internal/coremetrics/service.go new file mode 100644 index 000000000..3db113895 --- /dev/null +++ b/services/agents-api/internal/coremetrics/service.go @@ -0,0 +1,177 @@ +package coremetrics + +import ( + "context" + "errors" + "regexp" + "runtime" + "sync" + "time" +) + +const SampleInterval = 30 * time.Second +const retention = 7 * 24 * time.Hour +const sampleCapacity = int(retention/SampleInterval) + 2 + +var jobIDs = []string{"scheduler", "runtime_sampler", "history_cleanup", "audit_cleanup"} +var revisionPattern = regexp.MustCompile(`^[0-9a-f]{40}$`) + +type refusalSlot struct { + tick int64 + count int64 +} +type Service struct { + mu sync.Mutex + source Source + started time.Time + revision *string + now func() time.Time + samples [sampleCapacity]Sample + refusals [sampleCapacity]refusalSlot + latest Sample + jobs map[string]Job +} + +func New(started time.Time, revision string, source Source) *Service { + s := &Service{source: source, started: started.UTC(), now: time.Now, jobs: map[string]Job{}} + if revisionPattern.MatchString(revision) { + s.revision = &revision + } + for _, id := range jobIDs { + s.jobs[id] = Job{ID: id, Status: "unknown"} + } + return s +} +func slot(t time.Time) (int64, int) { + tick := t.Unix() / int64(SampleInterval/time.Second) + return tick, int(tick % int64(sampleCapacity)) +} + +// RecordUnavailable is called only when Core emits this exact rejection code. +// Fixed slots bound memory independently of request volume. +func (s *Service) RecordUnavailable() { + tick, i := slot(s.now()) + s.mu.Lock() + defer s.mu.Unlock() + if s.refusals[i].tick != tick { + s.refusals[i] = refusalSlot{tick: tick} + } + s.refusals[i].count++ +} +func (s *Service) ReportJob(id string, at time.Time, processed *int64, failed *int64, err error) { + status := "ok" + if err != nil || (failed != nil && *failed > 0) { + status = "failing" + } + s.mu.Lock() + defer s.mu.Unlock() + if _, ok := s.jobs[id]; !ok { + return + } + s.jobs[id] = Job{ID: id, Status: status, LastRunAt: ptr(at.UTC()), Processed: processed, Failed: failed} +} +func (s *Service) StopJob(id string) { + s.mu.Lock() + defer s.mu.Unlock() + if j, ok := s.jobs[id]; ok { + j.Status = "stopped" + s.jobs[id] = j + } +} +func (s *Service) Run(ctx context.Context) { + ticker := time.NewTicker(SampleInterval) + defer ticker.Stop() + for { + probe, cancel := context.WithTimeout(ctx, 5*time.Second) + sample := s.source.Sample(probe) + cancel() + sample.At = s.now().UTC() + s.record(sample) + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + } +} +func (s *Service) record(sample Sample) { + _, i := slot(sample.At) + s.mu.Lock() + defer s.mu.Unlock() + s.samples[i] = sample + s.latest = sample +} +func Window(now time.Time, name string) (Range, error) { + var duration, step time.Duration + switch name { + case "1h": + duration, step = time.Hour, time.Minute + case "6h": + duration, step = 6*time.Hour, 5*time.Minute + case "24h": + duration, step = 24*time.Hour, 15*time.Minute + case "7d": + duration, step = retention, 2*time.Hour + default: + return Range{}, errors.New("invalid Core metrics range") + } + end := now.UTC().Truncate(step) + return Range{Start: end.Add(-duration), End: end, ResolutionSeconds: int64(step / time.Second)}, nil +} +func (s *Service) Read(ctx context.Context, name string) (View, error) { + now := s.now().UTC() + window, err := Window(now, name) + if err != nil { + return View{}, err + } + live := s.source.Live() + view := View{Object: "core.metrics", Range: window, Service: ServiceState{Status: "running", Revision: s.revision, StartedAt: ptr(s.started), ExecutionOwner: live.ExecutionOwner}, + Execution: Execution{SlotsInUse: live.SlotsInUse, SlotsTotal: live.SlotsTotal, ConnectedDaemons: live.ConnectedDaemons}, Database: Database{Pool: live.Pool}, Jobs: make([]Job, 0, len(jobIDs))} + s.mu.Lock() + latest := s.latest + if latest.At.IsZero() || now.Sub(latest.At) > 2*SampleInterval || !latest.Healthy { + view.Service.Status = "degraded" + } + if !latest.At.IsZero() && now.Sub(latest.At) <= 2*SampleInterval { + view.Execution.QueuedTurns, view.Execution.WaitingForDaemon, view.Execution.InProgressTurns = latest.Queued, latest.WaitingForDaemon, latest.InProgress + view.Execution.OldestQueuedSeconds = latest.OldestQueuedSeconds + view.Database.SizeBytes = latest.DatabaseSize + if latest.Maintenance != nil && *latest.Maintenance && view.Service.Status == "running" { + view.Service.Status = "maintenance" + } + } + for _, id := range jobIDs { + j := s.jobs[id] + if id == "scheduler" && live.Scheduler.ID != "" { + j = live.Scheduler + } + view.Jobs = append(view.Jobs, j) + } + s.series(&view) + s.mu.Unlock() + query, cancel := context.WithTimeout(ctx, 3*time.Second) + defer cancel() + history, err := s.source.History(query, window.Start, window.End, time.Duration(window.ResolutionSeconds)*time.Second) + if err != nil { + view.Service.Status = "degraded" + } else { + view.Execution.Interrupted = ptr(history.Interrupted) + view.Execution.QueueWaitMS = history.QueueWaitMS + for i := range view.Execution.Series { + view.Execution.Series[i].QueueWaitP95MS = history.Buckets[view.Execution.Series[i].Start] + } + } + if live.ExecutionOwner == nil || (live.SlotsTotal != nil && *live.SlotsTotal > 0 && !*live.ExecutionOwner) { + view.Service.Status = "degraded" + } + for _, job := range view.Jobs { + if job.Status == "failing" { + view.Service.Status = "degraded" + } + } + var memory runtime.MemStats + runtime.ReadMemStats(&memory) + view.Process = Process{MemoryBytes: ptr(memory.Alloc), Goroutines: ptr(int64(runtime.NumGoroutine()))} + return view, nil +} +func ptr[T any](v T) *T { return &v } diff --git a/services/agents-api/internal/coremetrics/service_test.go b/services/agents-api/internal/coremetrics/service_test.go new file mode 100644 index 000000000..07281f506 --- /dev/null +++ b/services/agents-api/internal/coremetrics/service_test.go @@ -0,0 +1,161 @@ +package coremetrics + +import ( + "context" + "encoding/json" + "errors" + "strings" + "sync" + "testing" + "time" +) + +type fixtureSource struct { + history History + err error + sample Sample + live Live +} + +func (f *fixtureSource) Sample(context.Context) Sample { return f.sample } +func (f *fixtureSource) History(context.Context, time.Time, time.Time, time.Duration) (History, error) { + return f.history, f.err +} +func (f *fixtureSource) Live() Live { return f.live } +func fixtureService(t *testing.T) (*Service, *fixtureSource, time.Time) { + t.Helper() + now := time.Date(2026, 9, 25, 12, 0, 20, 0, time.UTC) + source := &fixtureSource{history: History{Buckets: map[time.Time]*float64{}}, live: Live{ExecutionOwner: ptr(true), SlotsInUse: ptr(int64(2)), SlotsTotal: ptr(int64(4))}} + service := New(now.Add(-2*time.Hour), strings.Repeat("a", 40), source) + service.now = func() time.Time { return now } + return service, source, now +} +func TestCompleteBucketsAndObservedPercentiles(t *testing.T) { + s, source, now := fixtureService(t) + end := now.Truncate(time.Minute) + for _, sample := range []Sample{ + {At: end.Add(-55 * time.Second), Queued: ptr(int64(3)), InProgress: ptr(int64(1)), PingMS: ptr(2.0), PoolInUse: ptr(int64(2)), Healthy: true}, + {At: end.Add(-15 * time.Second), Queued: ptr(int64(1)), InProgress: ptr(int64(4)), PingMS: ptr(6.0), PoolInUse: ptr(int64(3)), Healthy: true}, + {At: now, Queued: ptr(int64(99)), Healthy: true}, + } { + s.record(sample) + } + source.history = History{Interrupted: 2, QueueWaitMS: Latency{P50: ptr(100.0), P95: ptr(200.0)}, Buckets: map[time.Time]*float64{end.Add(-time.Minute): ptr(200.0)}} + s.now = func() time.Time { return end.Add(-10 * time.Second) } + s.RecordUnavailable() + s.now = func() time.Time { return now } + s.RecordUnavailable() + got, err := s.Read(t.Context(), "1h") + if err != nil { + t.Fatal(err) + } + if !got.Range.End.Equal(end) || len(got.Execution.Series) != 60 || *got.Execution.QueuedTurns != 99 || *got.Execution.Unavailable != 1 { + t.Fatal(got) + } + b := got.Execution.Series[59] + d := got.Database.Series[59] + if *b.Queued != 3 || *b.InProgress != 4 || *b.QueueWaitP95MS != 200 || *d.PoolInUse != 3 || *d.PingP95MS != 5.8 || *got.Database.PingMS.P50 != 4 { + t.Fatal(b, d, got.Database.PingMS) + } + if got.Execution.Series[0].Queued != nil || got.Database.Series[0].PingP95MS != nil { + t.Fatal("missing buckets became zero") + } +} +func TestUnknownAndFailureRemainNull(t *testing.T) { + s, source, now := fixtureService(t) + s.started = now.Add(-20 * time.Second) + source.err = errors.New("private database information") + s.record(Sample{At: now.Add(-10 * time.Second), PingMS: ptr(1.0), Queued: ptr(int64(0)), Healthy: false}) + got, err := s.Read(t.Context(), "1h") + if err != nil { + t.Fatal(err) + } + if got.Service.Status != "degraded" || got.Execution.Unavailable != nil || got.Execution.Interrupted != nil || got.Execution.QueueWaitMS.P50 != nil || got.Database.SizeBytes != nil { + t.Fatal(got) + } + raw, _ := json.Marshal(got) + if strings.Contains(string(raw), "private") || !strings.Contains(string(raw), `"unavailable":null`) { + t.Fatal(string(raw)) + } + for _, bucket := range got.Execution.Series { + if bucket.Queued != nil { + t.Fatal("pre-start bucket was filled", bucket) + } + } + s.started = now.Add(-time.Hour) + s.latest.At = now.Add(-3 * SampleInterval) + got, _ = s.Read(t.Context(), "1h") + if got.Execution.QueuedTurns != nil { + t.Fatal("stale gauge remained current") + } +} +func TestAllRangesAndBoundedRetention(t *testing.T) { + s, _, now := fixtureService(t) + for name, count := range map[string]int{"1h": 60, "6h": 72, "24h": 96, "7d": 84} { + got, err := s.Read(t.Context(), name) + if err != nil || len(got.Execution.Series) != count || len(got.Database.Series) != count { + t.Fatal(name, err) + } + } + if _, err := s.Read(t.Context(), "30d"); err == nil { + t.Fatal("unbounded range accepted") + } + // More than a retention window overwrites fixed slots, independent of load. + old := now.Add(-8 * 24 * time.Hour) + s.now = func() time.Time { return old } + s.RecordUnavailable() + s.record(Sample{At: old, Queued: ptr(int64(20))}) + s.now = func() time.Time { return now } + s.started = old + got, _ := s.Read(t.Context(), "7d") + if *got.Execution.Unavailable != 0 { + t.Fatal("expired refusal counted") + } + for _, v := range got.Execution.Series { + if v.Queued != nil { + t.Fatal("expired gauge counted") + } + } +} +func TestJobResultsAndConcurrentReads(t *testing.T) { + s, _, now := fixtureService(t) + s.record(Sample{At: now, Healthy: true, Maintenance: ptr(true)}) + s.ReportJob("audit_cleanup", now, ptr(int64(12)), ptr(int64(0)), nil) + got, _ := s.Read(t.Context(), "1h") + if got.Service.Status != "maintenance" || *got.Jobs[3].Processed != 12 { + t.Fatal(got.Service, got.Jobs) + } + s.ReportJob("audit_cleanup", now, ptr(int64(2)), ptr(int64(1)), errors.New("failure")) + got, _ = s.Read(t.Context(), "1h") + if got.Service.Status != "degraded" { + t.Fatal(got.Service) + } + s.StopJob("audit_cleanup") + got, _ = s.Read(t.Context(), "1h") + if got.Jobs[3].Status != "stopped" { + t.Fatal(got.Jobs) + } + var wg sync.WaitGroup + for range 4 { + wg.Add(1) + go func() { + defer wg.Done() + for range 20 { + s.RecordUnavailable() + s.ReportJob("history_cleanup", now, ptr(int64(0)), ptr(int64(0)), nil) + if _, err := s.Read(context.Background(), "1h"); err != nil { + t.Error(err) + } + } + }() + } + wg.Wait() +} +func TestRevisionMustBeCommit(t *testing.T) { + for _, revision := range []string{"", "unknown", "secret-value"} { + s := New(time.Now(), revision, &fixtureSource{}) + if s.revision != nil { + t.Fatal("unverified build revision exposed") + } + } +} diff --git a/services/agents-api/internal/coremetrics/types.go b/services/agents-api/internal/coremetrics/types.go new file mode 100644 index 000000000..c12de403c --- /dev/null +++ b/services/agents-api/internal/coremetrics/types.go @@ -0,0 +1,105 @@ +// Package coremetrics collects bounded operational measurements for administrators. +package coremetrics + +import ( + "context" + "time" +) + +type Latency struct { + P50 *float64 `json:"p50"` + P95 *float64 `json:"p95"` +} +type Range struct { + Start time.Time `json:"start"` + End time.Time `json:"end"` + ResolutionSeconds int64 `json:"resolution_seconds"` +} +type ServiceState struct { + Status string `json:"status"` + Revision *string `json:"revision"` + StartedAt *time.Time `json:"started_at"` + ExecutionOwner *bool `json:"execution_owner"` +} +type ExecutionBucket struct { + Start time.Time `json:"start"` + Queued *int64 `json:"queued"` + InProgress *int64 `json:"in_progress"` + QueueWaitP95MS *float64 `json:"queue_wait_p95_ms"` +} +type DatabaseBucket struct { + Start time.Time `json:"start"` + PingP95MS *float64 `json:"ping_p95_ms"` + PoolInUse *int64 `json:"pool_in_use"` +} +type Pool struct { + InUse *int64 `json:"in_use"` + Idle *int64 `json:"idle"` + Max *int64 `json:"max"` +} +type Execution struct { + SlotsInUse *int64 `json:"slots_in_use"` + SlotsTotal *int64 `json:"slots_total"` + QueuedTurns *int64 `json:"queued_turns"` + WaitingForDaemon *int64 `json:"waiting_for_daemon"` + InProgressTurns *int64 `json:"in_progress_turns"` + OldestQueuedSeconds *float64 `json:"oldest_queued_seconds"` + ConnectedDaemons *int64 `json:"connected_daemons"` + Interrupted *int64 `json:"interrupted"` + Unavailable *int64 `json:"unavailable"` + QueueWaitMS Latency `json:"queue_wait_ms"` + Series []ExecutionBucket `json:"series"` +} +type Database struct { + PingMS Latency `json:"ping_ms"` + Pool Pool `json:"pool"` + SizeBytes *int64 `json:"size_bytes"` + Series []DatabaseBucket `json:"series"` +} +type Job struct { + ID string `json:"id"` + Status string `json:"status"` + LastRunAt *time.Time `json:"last_run_at"` + Processed *int64 `json:"processed"` + Failed *int64 `json:"failed"` +} +type Process struct { + MemoryBytes *uint64 `json:"memory_bytes"` + Goroutines *int64 `json:"goroutines"` +} +type View struct { + Object string `json:"object"` + Range Range `json:"range"` + Service ServiceState `json:"service"` + Execution Execution `json:"execution"` + Database Database `json:"database"` + Jobs []Job `json:"jobs"` + Process Process `json:"process"` +} + +// Sample contains only measurements, never resource identities or error text. +type Sample struct { + At time.Time + Queued, WaitingForDaemon, InProgress *int64 + OldestQueuedSeconds *float64 + PingMS *float64 + PoolInUse, DatabaseSize *int64 + Maintenance *bool + Healthy bool +} +type Live struct { + SlotsInUse, SlotsTotal, ConnectedDaemons *int64 + ExecutionOwner *bool + Pool Pool + Scheduler Job +} +type History struct { + Interrupted int64 + QueueWaitMS Latency + Buckets map[time.Time]*float64 +} +type Source interface { + Sample(context.Context) Sample + History(context.Context, time.Time, time.Time, time.Duration) (History, error) + Live() Live +} diff --git a/services/agents-api/internal/runtimehistory/postgresreader/reader.go b/services/agents-api/internal/runtimehistory/postgresreader/reader.go index 7cced116c..54ffbe3b1 100644 --- a/services/agents-api/internal/runtimehistory/postgresreader/reader.go +++ b/services/agents-api/internal/runtimehistory/postgresreader/reader.go @@ -116,20 +116,22 @@ func (r *Reader) Query(ctx context.Context, query runtimehistory.Query) (runtime // Prune expires old telemetry even when no Runtime is being sampled. Each call // deletes at most 4,096 rows in small lock-skipping batches under one deadline. -func (r *Reader) Prune(ctx context.Context) error { +func (r *Reader) Prune(ctx context.Context) (int64, error) { ctx, cancel := context.WithTimeout(ctx, r.queryTimeout) defer cancel() before := r.now().Add(-retention).UnixNano() + var removed int64 for range 16 { count, err := r.store.PruneRuntimeHistorySamples(ctx, before) if err != nil { - return errors.New("prune Runtime history") + return removed, errors.New("prune Runtime history") } + removed += count if count < 256 { - return nil + return removed, nil } } - return nil + return removed, nil } func (r *Reader) validateQuery(query runtimehistory.Query) error { diff --git a/services/agents-api/internal/runtimehistory/postgresreader/reader_test.go b/services/agents-api/internal/runtimehistory/postgresreader/reader_test.go index c6f4b1171..a3a72c30f 100644 --- a/services/agents-api/internal/runtimehistory/postgresreader/reader_test.go +++ b/services/agents-api/internal/runtimehistory/postgresreader/reader_test.go @@ -166,11 +166,11 @@ func TestExporterPersistsOnlyPeriodicAndPreservesUnknownMetrics(t *testing.T) { func TestPruneHasBoundedWorkAndSanitizedFailure(t *testing.T) { s := &fakeStore{pruneCount: 256} r := testReader(t, s, time.Now()) - if err := r.Prune(t.Context()); err != nil || s.pruneCalls != 16 { + if count, err := r.Prune(t.Context()); err != nil || s.pruneCalls != 16 || count != 4096 { t.Fatal(err, s.pruneCalls) } s.err = errors.New("secret") - if err := r.Prune(t.Context()); err == nil || strings.Contains(err.Error(), "secret") { + if count, err := r.Prune(t.Context()); count != 0 || err == nil || strings.Contains(err.Error(), "secret") { t.Fatal(err) } } diff --git a/services/agents-api/internal/store/runtime_history_acceptance_test.go b/services/agents-api/internal/store/runtime_history_acceptance_test.go index 23cc3aa7f..dcedf2c99 100644 --- a/services/agents-api/internal/store/runtime_history_acceptance_test.go +++ b/services/agents-api/internal/store/runtime_history_acceptance_test.go @@ -184,14 +184,14 @@ func TestPostgresRuntimeHistoryDenseReadAndBoundedRetention(t *testing.T) { if err != nil || len(got.Coverage) != 0 { t.Fatal("expired samples visible", got, err) } - if err := reader.Prune(t.Context()); err != nil { + if _, err := reader.Prune(t.Context()); err != nil { t.Fatal(err) } var remaining int if err := pool.QueryRow(t.Context(), `SELECT count(*) FROM runtime_history_samples WHERE tenant_id=$1 AND resolved_at_ns<$2`, scope.TenantID, end.Add(-7*24*time.Hour).UnixNano()).Scan(&remaining); err != nil || remaining != 5 { t.Fatal("cleanup exceeded bounded batch", remaining, err) } - if err := reader.Prune(t.Context()); err != nil { + if _, err := reader.Prune(t.Context()); err != nil { t.Fatal(err) } if err := pool.QueryRow(t.Context(), `SELECT count(*) FROM runtime_history_samples WHERE tenant_id=$1 AND resolved_at_ns<$2`, scope.TenantID, end.Add(-7*24*time.Hour).UnixNano()).Scan(&remaining); err != nil || remaining != 0 { From 899664de44307d5b347c66fe58e85d3c117844b4 Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Fri, 25 Sep 2026 04:41:44 +0800 Subject: [PATCH 4/7] Retain complete seven-day metric buckets across alignment padding --- contracts/agents-api/core-metrics.md | 4 +-- .../internal/api/core_metrics_test.go | 3 ++- .../internal/coremetrics/service.go | 4 ++- .../internal/coremetrics/service_test.go | 25 +++++++++++++++++++ 4 files changed, 32 insertions(+), 4 deletions(-) diff --git a/contracts/agents-api/core-metrics.md b/contracts/agents-api/core-metrics.md index aa8d2f20b..4ac89bb19 100644 --- a/contracts/agents-api/core-metrics.md +++ b/contracts/agents-api/core-metrics.md @@ -34,8 +34,8 @@ that all intermediate peaks were captured. Missing observations and the process' partial first bucket stay null. Successful periodic ping samples produce linear interpolated p50/p95; there is no request-triggered ping. -A fixed-size in-process ring retains up to seven days of 30-second samples and -rejection counts. Restart loses those measurements: no synthetic backfill occurs. +A fixed-size in-process ring retains seven days of 30-second samples and +rejection counts, plus two hours of padding for complete bucket alignment. Restart loses those measurements: no synthetic backfill occurs. The `execution.unavailable` count is null if the requested interval starts before this process's observation began; an entirely observed interval with no rejections is zero. PostgreSQL Turn history remains queryable across process restarts. diff --git a/services/agents-api/internal/api/core_metrics_test.go b/services/agents-api/internal/api/core_metrics_test.go index bff441d14..4653a42b0 100644 --- a/services/agents-api/internal/api/core_metrics_test.go +++ b/services/agents-api/internal/api/core_metrics_test.go @@ -5,6 +5,7 @@ import ( "errors" "net/http" "net/http/httptest" + "strings" "testing" "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/device" @@ -60,7 +61,7 @@ func TestCoreMetricsAdministratorContract(t *testing.T) { t.Fatal(w.Code) } f.err = errors.New("private credentials") - if w := projectKeyHTTP(h, "GET", path, "admin", ""); w.Code != 503 || w.Body.String() == f.err.Error() { + if w := projectKeyHTTP(h, "GET", path, "admin", ""); w.Code != 503 || strings.Contains(w.Body.String(), f.err.Error()) { t.Fatal(w.Code, w.Body.String()) } } diff --git a/services/agents-api/internal/coremetrics/service.go b/services/agents-api/internal/coremetrics/service.go index 3db113895..d3f0b44bf 100644 --- a/services/agents-api/internal/coremetrics/service.go +++ b/services/agents-api/internal/coremetrics/service.go @@ -11,7 +11,9 @@ import ( const SampleInterval = 30 * time.Second const retention = 7 * 24 * time.Hour -const sampleCapacity = int(retention/SampleInterval) + 2 + +// The 7d window ends at a complete 2h bucket, so retain its leading padding too. +const sampleCapacity = int((retention+2*time.Hour)/SampleInterval) + 2 var jobIDs = []string{"scheduler", "runtime_sampler", "history_cleanup", "audit_cleanup"} var revisionPattern = regexp.MustCompile(`^[0-9a-f]{40}$`) diff --git a/services/agents-api/internal/coremetrics/service_test.go b/services/agents-api/internal/coremetrics/service_test.go index 07281f506..58f8b7717 100644 --- a/services/agents-api/internal/coremetrics/service_test.go +++ b/services/agents-api/internal/coremetrics/service_test.go @@ -159,3 +159,28 @@ func TestRevisionMustBeCommit(t *testing.T) { } } } + +func TestSevenDayAlignedWindowRetainsLeadingSamples(t *testing.T) { + s, _, now := fixtureService(t) + now = now.Add(time.Hour + 59*time.Minute) + window, err := Window(now, "7d") + if err != nil { + t.Fatal(err) + } + s.started = window.Start.Add(-time.Hour) + first := window.Start.Add(10 * time.Second) + s.now = func() time.Time { return first } + s.RecordUnavailable() + s.record(Sample{At: first, Queued: ptr(int64(7)), Healthy: true}) + // Populate every slot through the present, crossing the ordinary 7d cutoff. + for at := first.Add(SampleInterval); !at.After(now); at = at.Add(SampleInterval) { + s.record(Sample{At: at, Queued: ptr(int64(0)), Healthy: true}) + s.now = func() time.Time { return at } + s.RecordUnavailable() + } + s.now = func() time.Time { return now } + got, err := s.Read(t.Context(), "7d") + if err != nil || got.Execution.Unavailable == nil || *got.Execution.Unavailable != int64(retention/SampleInterval) || got.Execution.Series[0].Queued == nil || *got.Execution.Series[0].Queued != 7 { + t.Fatal("leading complete bucket was overwritten", err, got.Execution.Unavailable, got.Execution.Series[0]) + } +} From e662b867fe0211a6dd9148a66a0dfe4d7d80fcc5 Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Fri, 25 Sep 2026 04:42:41 +0800 Subject: [PATCH 5/7] Keep enabled sampler status unknown until its first result --- services/agents-api/cmd/server/main.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/services/agents-api/cmd/server/main.go b/services/agents-api/cmd/server/main.go index 9d63a7647..c86006bfb 100644 --- a/services/agents-api/cmd/server/main.go +++ b/services/agents-api/cmd/server/main.go @@ -93,7 +93,6 @@ func run() error { executionStore := store.NewWithCredentialCipherAndOAuthRefresh(pool, credentialKey, oauthClient) metricsSource := &coreMetricsSource{store: executionStore, pool: pool} metrics := coremetrics.New(processStartedAt, buildRevision, metricsSource) - metrics.StopJob("runtime_sampler") auth, err := api.NewDatabaseAuthenticator(executionStore) if err != nil { return err @@ -228,6 +227,9 @@ func run() error { options = append(options, api.WithHostedEnvironments()) } } + if history.SampleInterval == 0 { + metrics.StopJob("runtime_sampler") + } if history.SampleInterval > 0 { if worker == nil { return errors.New("Runtime history periodic sampling requires the execution worker") From d731cacae2ffbd82311971a7124caef6daa13e0b Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Fri, 25 Sep 2026 04:54:15 +0800 Subject: [PATCH 6/7] Declare nullable operational measurements in generated management schema --- .../agents-api/sandbox-manager.openapi.yaml | 28 ++++++++++ .../agents-api/internal/coremetrics/types.go | 56 +++++++++---------- 2 files changed, 56 insertions(+), 28 deletions(-) diff --git a/contracts/agents-api/sandbox-manager.openapi.yaml b/contracts/agents-api/sandbox-manager.openapi.yaml index 05baa293f..2fbec10e1 100644 --- a/contracts/agents-api/sandbox-manager.openapi.yaml +++ b/contracts/agents-api/sandbox-manager.openapi.yaml @@ -184,13 +184,16 @@ definitions: type: array size_bytes: type: integer + x-nullable: true type: object coremetrics.DatabaseBucket: properties: ping_p95_ms: type: number + x-nullable: true pool_in_use: type: integer + x-nullable: true start: type: string type: object @@ -198,37 +201,49 @@ definitions: properties: connected_daemons: type: integer + x-nullable: true in_progress_turns: type: integer + x-nullable: true interrupted: type: integer + x-nullable: true oldest_queued_seconds: type: number + x-nullable: true queue_wait_ms: $ref: '#/definitions/coremetrics.Latency' queued_turns: type: integer + x-nullable: true series: items: $ref: '#/definitions/coremetrics.ExecutionBucket' type: array slots_in_use: type: integer + x-nullable: true slots_total: type: integer + x-nullable: true unavailable: type: integer + x-nullable: true waiting_for_daemon: type: integer + x-nullable: true type: object coremetrics.ExecutionBucket: properties: in_progress: type: integer + x-nullable: true queue_wait_p95_ms: type: number + x-nullable: true queued: type: integer + x-nullable: true start: type: string type: object @@ -236,12 +251,15 @@ definitions: properties: failed: type: integer + x-nullable: true id: type: string last_run_at: type: string + x-nullable: true processed: type: integer + x-nullable: true status: type: string type: object @@ -249,24 +267,31 @@ definitions: properties: p50: type: number + x-nullable: true p95: type: number + x-nullable: true type: object coremetrics.Pool: properties: idle: type: integer + x-nullable: true in_use: type: integer + x-nullable: true max: type: integer + x-nullable: true type: object coremetrics.Process: properties: goroutines: type: integer + x-nullable: true memory_bytes: type: integer + x-nullable: true type: object coremetrics.Range: properties: @@ -281,10 +306,13 @@ definitions: properties: execution_owner: type: boolean + x-nullable: true revision: type: string + x-nullable: true started_at: type: string + x-nullable: true status: type: string type: object diff --git a/services/agents-api/internal/coremetrics/types.go b/services/agents-api/internal/coremetrics/types.go index c12de403c..2a926604c 100644 --- a/services/agents-api/internal/coremetrics/types.go +++ b/services/agents-api/internal/coremetrics/types.go @@ -7,8 +7,8 @@ import ( ) type Latency struct { - P50 *float64 `json:"p50"` - P95 *float64 `json:"p95"` + P50 *float64 `json:"p50" extensions:"x-nullable"` + P95 *float64 `json:"p95" extensions:"x-nullable"` } type Range struct { Start time.Time `json:"start"` @@ -17,55 +17,55 @@ type Range struct { } type ServiceState struct { Status string `json:"status"` - Revision *string `json:"revision"` - StartedAt *time.Time `json:"started_at"` - ExecutionOwner *bool `json:"execution_owner"` + Revision *string `json:"revision" extensions:"x-nullable"` + StartedAt *time.Time `json:"started_at" extensions:"x-nullable"` + ExecutionOwner *bool `json:"execution_owner" extensions:"x-nullable"` } type ExecutionBucket struct { Start time.Time `json:"start"` - Queued *int64 `json:"queued"` - InProgress *int64 `json:"in_progress"` - QueueWaitP95MS *float64 `json:"queue_wait_p95_ms"` + Queued *int64 `json:"queued" extensions:"x-nullable"` + InProgress *int64 `json:"in_progress" extensions:"x-nullable"` + QueueWaitP95MS *float64 `json:"queue_wait_p95_ms" extensions:"x-nullable"` } type DatabaseBucket struct { Start time.Time `json:"start"` - PingP95MS *float64 `json:"ping_p95_ms"` - PoolInUse *int64 `json:"pool_in_use"` + PingP95MS *float64 `json:"ping_p95_ms" extensions:"x-nullable"` + PoolInUse *int64 `json:"pool_in_use" extensions:"x-nullable"` } type Pool struct { - InUse *int64 `json:"in_use"` - Idle *int64 `json:"idle"` - Max *int64 `json:"max"` + InUse *int64 `json:"in_use" extensions:"x-nullable"` + Idle *int64 `json:"idle" extensions:"x-nullable"` + Max *int64 `json:"max" extensions:"x-nullable"` } type Execution struct { - SlotsInUse *int64 `json:"slots_in_use"` - SlotsTotal *int64 `json:"slots_total"` - QueuedTurns *int64 `json:"queued_turns"` - WaitingForDaemon *int64 `json:"waiting_for_daemon"` - InProgressTurns *int64 `json:"in_progress_turns"` - OldestQueuedSeconds *float64 `json:"oldest_queued_seconds"` - ConnectedDaemons *int64 `json:"connected_daemons"` - Interrupted *int64 `json:"interrupted"` - Unavailable *int64 `json:"unavailable"` + SlotsInUse *int64 `json:"slots_in_use" extensions:"x-nullable"` + SlotsTotal *int64 `json:"slots_total" extensions:"x-nullable"` + QueuedTurns *int64 `json:"queued_turns" extensions:"x-nullable"` + WaitingForDaemon *int64 `json:"waiting_for_daemon" extensions:"x-nullable"` + InProgressTurns *int64 `json:"in_progress_turns" extensions:"x-nullable"` + OldestQueuedSeconds *float64 `json:"oldest_queued_seconds" extensions:"x-nullable"` + ConnectedDaemons *int64 `json:"connected_daemons" extensions:"x-nullable"` + Interrupted *int64 `json:"interrupted" extensions:"x-nullable"` + Unavailable *int64 `json:"unavailable" extensions:"x-nullable"` QueueWaitMS Latency `json:"queue_wait_ms"` Series []ExecutionBucket `json:"series"` } type Database struct { PingMS Latency `json:"ping_ms"` Pool Pool `json:"pool"` - SizeBytes *int64 `json:"size_bytes"` + SizeBytes *int64 `json:"size_bytes" extensions:"x-nullable"` Series []DatabaseBucket `json:"series"` } type Job struct { ID string `json:"id"` Status string `json:"status"` - LastRunAt *time.Time `json:"last_run_at"` - Processed *int64 `json:"processed"` - Failed *int64 `json:"failed"` + LastRunAt *time.Time `json:"last_run_at" extensions:"x-nullable"` + Processed *int64 `json:"processed" extensions:"x-nullable"` + Failed *int64 `json:"failed" extensions:"x-nullable"` } type Process struct { - MemoryBytes *uint64 `json:"memory_bytes"` - Goroutines *int64 `json:"goroutines"` + MemoryBytes *uint64 `json:"memory_bytes" extensions:"x-nullable"` + Goroutines *int64 `json:"goroutines" extensions:"x-nullable"` } type View struct { Object string `json:"object"` From 622693f37c210d7bd4692fad36195571fa8afdee Mon Sep 17 00:00:00 2001 From: saladday <1203511142@qq.com> Date: Fri, 25 Sep 2026 04:58:02 +0800 Subject: [PATCH 7/7] Retain every completed metrics sample despite timing jitter --- .../internal/coremetrics/service.go | 24 ++++++++++--------- .../internal/coremetrics/service_test.go | 5 ++-- 2 files changed, 16 insertions(+), 13 deletions(-) diff --git a/services/agents-api/internal/coremetrics/service.go b/services/agents-api/internal/coremetrics/service.go index d3f0b44bf..646451d09 100644 --- a/services/agents-api/internal/coremetrics/service.go +++ b/services/agents-api/internal/coremetrics/service.go @@ -23,15 +23,16 @@ type refusalSlot struct { count int64 } type Service struct { - mu sync.Mutex - source Source - started time.Time - revision *string - now func() time.Time - samples [sampleCapacity]Sample - refusals [sampleCapacity]refusalSlot - latest Sample - jobs map[string]Job + mu sync.Mutex + source Source + started time.Time + revision *string + now func() time.Time + samples [sampleCapacity]Sample + nextSample int + refusals [sampleCapacity]refusalSlot + latest Sample + jobs map[string]Job } func New(started time.Time, revision string, source Source) *Service { @@ -97,10 +98,11 @@ func (s *Service) Run(ctx context.Context) { } } func (s *Service) record(sample Sample) { - _, i := slot(sample.At) s.mu.Lock() defer s.mu.Unlock() - s.samples[i] = sample + // Consecutive probes can finish within the same wall-clock slot. Keep both. + s.samples[s.nextSample] = sample + s.nextSample = (s.nextSample + 1) % len(s.samples) s.latest = sample } func Window(now time.Time, name string) (Range, error) { diff --git a/services/agents-api/internal/coremetrics/service_test.go b/services/agents-api/internal/coremetrics/service_test.go index 58f8b7717..fb6daf187 100644 --- a/services/agents-api/internal/coremetrics/service_test.go +++ b/services/agents-api/internal/coremetrics/service_test.go @@ -33,9 +33,10 @@ func fixtureService(t *testing.T) (*Service, *fixtureSource, time.Time) { func TestCompleteBucketsAndObservedPercentiles(t *testing.T) { s, source, now := fixtureService(t) end := now.Truncate(time.Minute) + // The first two probes finish in the same 30s slot after differing I/O delays. for _, sample := range []Sample{ - {At: end.Add(-55 * time.Second), Queued: ptr(int64(3)), InProgress: ptr(int64(1)), PingMS: ptr(2.0), PoolInUse: ptr(int64(2)), Healthy: true}, - {At: end.Add(-15 * time.Second), Queued: ptr(int64(1)), InProgress: ptr(int64(4)), PingMS: ptr(6.0), PoolInUse: ptr(int64(3)), Healthy: true}, + {At: end.Add(-27 * time.Second), Queued: ptr(int64(3)), InProgress: ptr(int64(1)), PingMS: ptr(2.0), PoolInUse: ptr(int64(2)), Healthy: true}, + {At: end.Add(-time.Second), Queued: ptr(int64(1)), InProgress: ptr(int64(4)), PingMS: ptr(6.0), PoolInUse: ptr(int64(3)), Healthy: true}, {At: now, Queued: ptr(int64(99)), Healthy: true}, } { s.record(sample)