mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: change agent to use v2 API for reporting stats (#12024)
Modifies the agent to use the v2 API to report its statistics, using the `statsReporter` subcomponent.
This commit is contained in:
@@ -24,7 +24,6 @@ import (
|
||||
"github.com/coder/coder/v2/agent/proto"
|
||||
"github.com/coder/coder/v2/codersdk"
|
||||
drpcsdk "github.com/coder/coder/v2/codersdk/drpc"
|
||||
"github.com/coder/retry"
|
||||
)
|
||||
|
||||
// ExternalLogSourceID is the statically-defined ID of a log-source that
|
||||
@@ -390,61 +389,6 @@ func (c *Client) AuthAzureInstanceIdentity(ctx context.Context) (AuthenticateRes
|
||||
return resp, json.NewDecoder(res.Body).Decode(&resp)
|
||||
}
|
||||
|
||||
// ReportStats begins a stat streaming connection with the Coder server.
|
||||
// It is resilient to network failures and intermittent coderd issues.
|
||||
func (c *Client) ReportStats(ctx context.Context, log slog.Logger, statsChan <-chan *Stats, setInterval func(time.Duration)) (io.Closer, error) {
|
||||
var interval time.Duration
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
exited := make(chan struct{})
|
||||
|
||||
postStat := func(stat *Stats) {
|
||||
var nextInterval time.Duration
|
||||
for r := retry.New(100*time.Millisecond, time.Minute); r.Wait(ctx); {
|
||||
resp, err := c.PostStats(ctx, stat)
|
||||
if err != nil {
|
||||
if !xerrors.Is(err, context.Canceled) {
|
||||
log.Error(ctx, "report stats", slog.Error(err))
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
nextInterval = resp.ReportInterval
|
||||
break
|
||||
}
|
||||
|
||||
if nextInterval != 0 && interval != nextInterval {
|
||||
setInterval(nextInterval)
|
||||
}
|
||||
interval = nextInterval
|
||||
}
|
||||
|
||||
// Send an empty stat to get the interval.
|
||||
postStat(&Stats{})
|
||||
|
||||
go func() {
|
||||
defer close(exited)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case stat, ok := <-statsChan:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
postStat(stat)
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return closeFunc(func() error {
|
||||
cancel()
|
||||
<-exited
|
||||
return nil
|
||||
}), nil
|
||||
}
|
||||
|
||||
// Stats records the Agent's network connection statistics for use in
|
||||
// user-facing metrics and debugging.
|
||||
type Stats struct {
|
||||
@@ -509,6 +453,9 @@ type StatsResponse struct {
|
||||
ReportInterval time.Duration `json:"report_interval"`
|
||||
}
|
||||
|
||||
// PostStats sends agent stats to the coder server
|
||||
//
|
||||
// Deprecated: uses agent API v1 endpoint
|
||||
func (c *Client) PostStats(ctx context.Context, stats *Stats) (StatsResponse, error) {
|
||||
res, err := c.SDK.Request(ctx, http.MethodPost, "/api/v2/workspaceagents/me/report-stats", stats)
|
||||
if err != nil {
|
||||
@@ -649,12 +596,6 @@ func (c *Client) ExternalAuth(ctx context.Context, req ExternalAuthRequest) (Ext
|
||||
return authResp, json.NewDecoder(res.Body).Decode(&authResp)
|
||||
}
|
||||
|
||||
type closeFunc func() error
|
||||
|
||||
func (c closeFunc) Close() error {
|
||||
return c()
|
||||
}
|
||||
|
||||
// wsNetConn wraps net.Conn created by websocket.NetConn(). Cancel func
|
||||
// is called if a read or write error is encountered.
|
||||
type wsNetConn struct {
|
||||
|
||||
@@ -1,22 +1,13 @@
|
||||
package codersdk_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"tailscale.com/tailcfg"
|
||||
|
||||
"cdr.dev/slog/sloggers/slogtest"
|
||||
"github.com/coder/coder/v2/coderd/httpapi"
|
||||
"github.com/coder/coder/v2/codersdk/agentsdk"
|
||||
"github.com/coder/coder/v2/testutil"
|
||||
)
|
||||
|
||||
func TestWorkspaceRewriteDERPMap(t *testing.T) {
|
||||
@@ -46,45 +37,3 @@ func TestWorkspaceRewriteDERPMap(t *testing.T) {
|
||||
require.Equal(t, "coconuts.org", node.HostName)
|
||||
require.Equal(t, 44558, node.DERPPort)
|
||||
}
|
||||
|
||||
func TestAgentReportStats(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var (
|
||||
numReports atomic.Int64
|
||||
numIntervalCalls atomic.Int64
|
||||
wantInterval = 5 * time.Millisecond
|
||||
)
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
numReports.Add(1)
|
||||
httpapi.Write(context.Background(), w, http.StatusOK, agentsdk.StatsResponse{
|
||||
ReportInterval: wantInterval,
|
||||
})
|
||||
}))
|
||||
parsed, err := url.Parse(srv.URL)
|
||||
require.NoError(t, err)
|
||||
client := agentsdk.New(parsed)
|
||||
|
||||
assertStatInterval := func(interval time.Duration) {
|
||||
numIntervalCalls.Add(1)
|
||||
assert.Equal(t, wantInterval, interval)
|
||||
}
|
||||
|
||||
chanLen := 3
|
||||
statCh := make(chan *agentsdk.Stats, chanLen)
|
||||
for i := 0; i < chanLen; i++ {
|
||||
statCh <- &agentsdk.Stats{ConnectionsByProto: map[string]int64{}}
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
closeStream, err := client.ReportStats(ctx, slogtest.Make(t, nil), statCh, assertStatInterval)
|
||||
require.NoError(t, err)
|
||||
defer closeStream.Close()
|
||||
|
||||
require.Eventually(t,
|
||||
func() bool { return numReports.Load() >= 3 },
|
||||
testutil.WaitMedium, testutil.IntervalFast,
|
||||
)
|
||||
closeStream.Close()
|
||||
require.Equal(t, int64(1), numIntervalCalls.Load())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user