Compare commits
46
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6fcfe31040 | ||
|
|
bc6c5ed3e8 | ||
|
|
c250298318 | ||
|
|
14887afcce | ||
|
|
aeda0a5af4 | ||
|
|
c841c96f1c | ||
|
|
145b7e382f | ||
|
|
0015641f99 | ||
|
|
1208e1216e | ||
|
|
ebe67dba45 | ||
|
|
5cb859334b | ||
|
|
3fe477d93e | ||
|
|
f3ab19d2b8 | ||
|
|
40342b1e2a | ||
|
|
d9d8d097da | ||
|
|
e346584c2d | ||
|
|
be0376b5b5 | ||
|
|
7e0eae99a6 | ||
|
|
01adae723f | ||
|
|
e37307f46f | ||
|
|
23bfd4848d | ||
|
|
a2a8af97b6 | ||
|
|
60ce378d14 | ||
|
|
dde6b4631c | ||
|
|
fe39658434 | ||
|
|
cec89c5b7b | ||
|
|
7442885c99 | ||
|
|
6d1c317b90 | ||
|
|
45800be002 | ||
|
|
16a9c915cb | ||
|
|
c2f0a53402 | ||
|
|
5f370484a7 | ||
|
|
8104f571eb | ||
|
|
4cc97e220a | ||
|
|
8155f87a58 | ||
|
|
83fb2e2837 | ||
|
|
1f22d2441c | ||
|
|
e9e5ca3282 | ||
|
|
51aa04a4b2 | ||
|
|
39ffa2e560 | ||
|
|
b1dcbc44e6 | ||
|
|
287d3410da | ||
|
|
c549549283 | ||
|
|
a2aa9f9ea9 | ||
|
|
dd4c9f2481 | ||
|
|
8b43e9e5c9 |
No files matched your search
@@ -65,4 +65,4 @@ jobs:
|
|||||||
go-version-file: "go.mod"
|
go-version-file: "go.mod"
|
||||||
cache: true
|
cache: true
|
||||||
- name: Run tests
|
- name: Run tests
|
||||||
run: go test ./... -race
|
run: go test ./...
|
||||||
@@ -1,6 +1,6 @@
|
|||||||
MIT License
|
MIT License
|
||||||
|
|
||||||
Copyright GitHub, Inc.
|
Copyright (c) 2025 GitHub
|
||||||
|
|
||||||
Permission is hereby granted, free of charge, to any person obtaining a copy
|
Permission is hereby granted, free of charge, to any person obtaining a copy
|
||||||
of this software and associated documentation files (the "Software"), to deal
|
of this software and associated documentation files (the "Software"), to deal
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
# GitHub Actions Runner Scale Set Client (Public Preview)
|
# GitHub Actions Runner Scale Set Client (Private Preview)
|
||||||
|
|
||||||
> Status: **Public Preview** – While the API is stable, interfaces and examples in this repository may change.
|
> Status: **Private Preview** – While the API is stable, interfaces and examples in this repository may change.
|
||||||
|
|
||||||
This repository provides a standalone Go client for the GitHub Actions **Runner Scale Set** APIs. It is extracted from the `actions-runner-controller` project so that platform teams, integrators, and infrastructure providers can build **their own custom autoscaling solutions** for GitHub Actions runners.
|
This repository provides a standalone Go client for the GitHub Actions **Runner Scale Set** APIs. It is extracted from the `actions-runner-controller` project so that platform teams, integrators, and infrastructure providers can build **their own custom autoscaling solutions** for GitHub Actions runners.
|
||||||
|
|
||||||
@@ -12,7 +12,7 @@ You do *not* need to adopt the full controller (and Kubernetes) to take advantag
|
|||||||
|
|
||||||
A runner scale set is a group of self-hosted runners that autoscales based on workflow demand. Here's how it works:
|
A runner scale set is a group of self-hosted runners that autoscales based on workflow demand. Here's how it works:
|
||||||
|
|
||||||
1. **Registration**: You create a scale set with a name, which also serves as the label workflows use to target it (e.g., `runs-on: my-scale-set`). Multiple labels can be assigned per scale set. Like regular self-hosted runners, scale sets can be registered at the repository, organization, or enterprise level.
|
1. **Registration**: You create a scale set with a name, which also serves as the label workflows use to target it (e.g., `runs-on: my-scale-set`). Like regular self-hosted runners, scale sets can be registered at the repository, organization, or enterprise level.
|
||||||
2. **Polling**: Your scale set client continuously polls the API, reporting its maximum capacity (how many runners it can produce).
|
2. **Polling**: Your scale set client continuously polls the API, reporting its maximum capacity (how many runners it can produce).
|
||||||
3. **Job matching**: GitHub matches jobs to your scale set based on the label and runner group policies, just like regular self-hosted runners.
|
3. **Job matching**: GitHub matches jobs to your scale set based on the label and runner group policies, just like regular self-hosted runners.
|
||||||
4. **Scaling signal**: The API responds with how many runners your scale set needs online (`statistics.TotalAssignedJobs`).
|
4. **Scaling signal**: The API responds with how many runners your scale set needs online (`statistics.TotalAssignedJobs`).
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
|
"log/slog"
|
||||||
"maps"
|
"maps"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/url"
|
"net/url"
|
||||||
@@ -126,9 +127,9 @@ type ProxyFunc func(req *http.Request) (*url.URL, error)
|
|||||||
type SystemInfo struct {
|
type SystemInfo struct {
|
||||||
// System is the name of the scale set implementation
|
// System is the name of the scale set implementation
|
||||||
System string `json:"system"`
|
System string `json:"system"`
|
||||||
// Version is the version of the client
|
// Version is the version of the controller
|
||||||
Version string `json:"version"`
|
Version string `json:"version"`
|
||||||
// CommitSHA is the git commit SHA of the client
|
// CommitSHA is the git commit SHA of the controller
|
||||||
CommitSHA string `json:"commit_sha"`
|
CommitSHA string `json:"commit_sha"`
|
||||||
// ScaleSetID is the ID of the scale set
|
// ScaleSetID is the ID of the scale set
|
||||||
ScaleSetID int `json:"scale_set_id"`
|
ScaleSetID int `json:"scale_set_id"`
|
||||||
@@ -222,7 +223,6 @@ type userAgent struct {
|
|||||||
SystemInfo
|
SystemInfo
|
||||||
BuildVersion string `json:"build_version"`
|
BuildVersion string `json:"build_version"`
|
||||||
BuildCommitSHA string `json:"build_commit_sha"`
|
BuildCommitSHA string `json:"build_commit_sha"`
|
||||||
Kind string `json:"kind"`
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Client) newGitHubAPIRequest(ctx context.Context, method, path string, body io.Reader) (*http.Request, error) {
|
func (c *Client) newGitHubAPIRequest(ctx context.Context, method, path string, body io.Reader) (*http.Request, error) {
|
||||||
@@ -546,6 +546,7 @@ func parseRunnerScaleSetMessageResponse(respBody io.Reader) (*RunnerScaleSetMess
|
|||||||
// It exposes client options that could be overwritten, providing ability to specify different retry policies or TLS settings, proxy, etc.
|
// It exposes client options that could be overwritten, providing ability to specify different retry policies or TLS settings, proxy, etc.
|
||||||
func (c *Client) MessageSessionClient(ctx context.Context, runnerScaleSetID int, owner string, options ...HTTPOption) (*MessageSessionClient, error) {
|
func (c *Client) MessageSessionClient(ctx context.Context, runnerScaleSetID int, owner string, options ...HTTPOption) (*MessageSessionClient, error) {
|
||||||
c.mu.Lock()
|
c.mu.Lock()
|
||||||
|
defer c.mu.Unlock()
|
||||||
|
|
||||||
// Copy original options
|
// Copy original options
|
||||||
httpClientOption := c.httpClientOption
|
httpClientOption := c.httpClientOption
|
||||||
@@ -566,8 +567,6 @@ func (c *Client) MessageSessionClient(ctx context.Context, runnerScaleSetID int,
|
|||||||
scaleSetID: runnerScaleSetID,
|
scaleSetID: runnerScaleSetID,
|
||||||
session: nil,
|
session: nil,
|
||||||
}
|
}
|
||||||
// Unlock the client to allow createMessageSession to call public methods that require locking
|
|
||||||
c.mu.Unlock()
|
|
||||||
|
|
||||||
if err := client.createMessageSession(ctx); err != nil {
|
if err := client.createMessageSession(ctx); err != nil {
|
||||||
return nil, fmt.Errorf("failed to create message session: %w", err)
|
return nil, fmt.Errorf("failed to create message session: %w", err)
|
||||||
@@ -837,6 +836,8 @@ func (c *Client) getActionsServiceAdminConnection(ctx context.Context, rt *regis
|
|||||||
return nil, fmt.Errorf("failed to get actions service admin connection: %w", err)
|
return nil, fmt.Errorf("failed to get actions service admin connection: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
slog.Info("got admin connection", *adminConnection.ActionsServiceURL)
|
||||||
|
|
||||||
return adminConnection, nil
|
return adminConnection, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+11
-55
@@ -75,20 +75,12 @@ func sendRequest(c *http.Client, req *http.Request) (*http.Response, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type httpClientOption struct {
|
type httpClientOption struct {
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
retryMax int
|
||||||
// Options for built-in retryable HTTP client.
|
retryWaitMax time.Duration
|
||||||
// Ignored if a custom retryable HTTP client is provided via WithRetryableHTTPClint.
|
|
||||||
retryMax int
|
|
||||||
retryWaitMax time.Duration
|
|
||||||
|
|
||||||
// fields added to the transport if specified
|
|
||||||
rootCAs *x509.CertPool
|
rootCAs *x509.CertPool
|
||||||
tlsInsecureSkipVerify bool
|
tlsInsecureSkipVerify bool
|
||||||
proxyFunc ProxyFunc
|
proxyFunc ProxyFunc
|
||||||
timeout time.Duration
|
|
||||||
|
|
||||||
retryableHTTPClient *retryablehttp.Client
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (o *httpClientOption) defaults() {
|
func (o *httpClientOption) defaults() {
|
||||||
@@ -101,28 +93,14 @@ func (o *httpClientOption) defaults() {
|
|||||||
if o.retryWaitMax == 0 {
|
if o.retryWaitMax == 0 {
|
||||||
o.retryWaitMax = 30 * time.Second
|
o.retryWaitMax = 30 * time.Second
|
||||||
}
|
}
|
||||||
if o.timeout == 0 {
|
|
||||||
o.timeout = 5 * time.Minute
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (o *httpClientOption) newRetryableHTTPClient() (*retryablehttp.Client, error) {
|
func (o *httpClientOption) newRetryableHTTPClient() (*retryablehttp.Client, error) {
|
||||||
var retryClient *retryablehttp.Client
|
retryClient := retryablehttp.NewClient()
|
||||||
if o.retryableHTTPClient != nil {
|
retryClient.Logger = o.logger
|
||||||
retryClient = o.retryableHTTPClient
|
retryClient.RetryMax = o.retryMax
|
||||||
} else {
|
retryClient.RetryWaitMax = o.retryWaitMax
|
||||||
retryClient = retryablehttp.NewClient()
|
retryClient.HTTPClient.Timeout = 5 * time.Minute // timeout must be > 1m to accomodate long polling
|
||||||
retryClient.RetryMax = o.retryMax
|
|
||||||
retryClient.RetryWaitMax = o.retryWaitMax
|
|
||||||
}
|
|
||||||
|
|
||||||
if retryClient.HTTPClient.Timeout == 0 {
|
|
||||||
retryClient.HTTPClient.Timeout = o.timeout
|
|
||||||
}
|
|
||||||
|
|
||||||
if retryClient.Logger == nil {
|
|
||||||
retryClient.Logger = o.logger
|
|
||||||
}
|
|
||||||
|
|
||||||
transport, ok := retryClient.HTTPClient.Transport.(*http.Transport)
|
transport, ok := retryClient.HTTPClient.Transport.(*http.Transport)
|
||||||
if !ok {
|
if !ok {
|
||||||
@@ -142,9 +120,7 @@ func (o *httpClientOption) newRetryableHTTPClient() (*retryablehttp.Client, erro
|
|||||||
transport.TLSClientConfig.InsecureSkipVerify = true
|
transport.TLSClientConfig.InsecureSkipVerify = true
|
||||||
}
|
}
|
||||||
|
|
||||||
if o.proxyFunc != nil {
|
transport.Proxy = o.proxyFunc
|
||||||
transport.Proxy = o.proxyFunc
|
|
||||||
}
|
|
||||||
|
|
||||||
retryClient.HTTPClient.Transport = transport
|
retryClient.HTTPClient.Transport = transport
|
||||||
|
|
||||||
@@ -161,7 +137,6 @@ func (c *commonClient) setUserAgent() {
|
|||||||
SystemInfo: c.systemInfo,
|
SystemInfo: c.systemInfo,
|
||||||
BuildVersion: buildInfo.version,
|
BuildVersion: buildInfo.version,
|
||||||
BuildCommitSHA: buildInfo.commitSHA,
|
BuildCommitSHA: buildInfo.commitSHA,
|
||||||
Kind: "scaleset",
|
|
||||||
})
|
})
|
||||||
c.userAgent = string(b)
|
c.userAgent = string(b)
|
||||||
}
|
}
|
||||||
@@ -169,22 +144,10 @@ func (c *commonClient) setUserAgent() {
|
|||||||
// HTTPOption defines a functional option for configuring the Client.
|
// HTTPOption defines a functional option for configuring the Client.
|
||||||
type HTTPOption func(*httpClientOption)
|
type HTTPOption func(*httpClientOption)
|
||||||
|
|
||||||
// WithRetryableHTTPClint allows users to provide a custom retryable HTTP client.
|
|
||||||
// If not set, a default client will be used with the specified retry and timeout settings.
|
|
||||||
func WithRetryableHTTPClint(client *retryablehttp.Client) HTTPOption {
|
|
||||||
return func(c *httpClientOption) {
|
|
||||||
c.retryableHTTPClient = client
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// WithLogger sets a custom logger for the Client.
|
// WithLogger sets a custom logger for the Client.
|
||||||
// If nil is passed, a discard logger will be used.
|
func WithLogger(logger slog.Logger) HTTPOption {
|
||||||
func WithLogger(logger *slog.Logger) HTTPOption {
|
|
||||||
return func(c *httpClientOption) {
|
return func(c *httpClientOption) {
|
||||||
if logger == nil {
|
c.logger = &logger
|
||||||
logger = slog.New(slog.DiscardHandler)
|
|
||||||
}
|
|
||||||
c.logger = logger
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -222,10 +185,3 @@ func WithProxy(proxyFunc ProxyFunc) HTTPOption {
|
|||||||
c.proxyFunc = proxyFunc
|
c.proxyFunc = proxyFunc
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// WithTimeout sets a timeout for the Client.
|
|
||||||
func WithTimeout(duration time.Duration) HTTPOption {
|
|
||||||
return func(c *httpClientOption) {
|
|
||||||
c.timeout = duration
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -7,10 +7,8 @@ import (
|
|||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"net/url"
|
"net/url"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/actions/scaleset/internal/testserver"
|
"github.com/actions/scaleset/internal/testserver"
|
||||||
"github.com/hashicorp/go-retryablehttp"
|
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
"golang.org/x/net/http/httpproxy"
|
"golang.org/x/net/http/httpproxy"
|
||||||
@@ -126,7 +124,6 @@ func TestUserAgent(t *testing.T) {
|
|||||||
SystemInfo: testSystemInfo,
|
SystemInfo: testSystemInfo,
|
||||||
BuildCommitSHA: sha,
|
BuildCommitSHA: sha,
|
||||||
BuildVersion: version,
|
BuildVersion: version,
|
||||||
Kind: "scaleset",
|
|
||||||
}
|
}
|
||||||
b, err := json.Marshal(wantInfo)
|
b, err := json.Marshal(wantInfo)
|
||||||
require.NoError(t, err, "failed to marshal expected user agent")
|
require.NoError(t, err, "failed to marshal expected user agent")
|
||||||
@@ -147,7 +144,6 @@ func TestUserAgent(t *testing.T) {
|
|||||||
SystemInfo: userAgentInfo,
|
SystemInfo: userAgentInfo,
|
||||||
BuildCommitSHA: sha,
|
BuildCommitSHA: sha,
|
||||||
BuildVersion: version,
|
BuildVersion: version,
|
||||||
Kind: "scaleset",
|
|
||||||
}
|
}
|
||||||
b, err = json.Marshal(wantInfo)
|
b, err = json.Marshal(wantInfo)
|
||||||
require.NoError(t, err, "failed to marshal expected user agent after SetSystemInfo")
|
require.NoError(t, err, "failed to marshal expected user agent after SetSystemInfo")
|
||||||
@@ -155,88 +151,3 @@ func TestUserAgent(t *testing.T) {
|
|||||||
|
|
||||||
assert.Equal(t, want, got)
|
assert.Equal(t, want, got)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestWithRetryableHTTPClient verifies that a custom retryable HTTP client
|
|
||||||
// provided via WithRetryableHTTPClient is actually used instead of the built-in one
|
|
||||||
func TestWithRetryableHTTPClient(t *testing.T) {
|
|
||||||
t.Run("uses custom retryable client instead of built-in", func(t *testing.T) {
|
|
||||||
attemptCount := 0
|
|
||||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
||||||
attemptCount++
|
|
||||||
if attemptCount == 1 {
|
|
||||||
w.WriteHeader(http.StatusServiceUnavailable)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
w.WriteHeader(http.StatusOK)
|
|
||||||
w.Write([]byte(`{"result": "success"}`))
|
|
||||||
}))
|
|
||||||
defer server.Close()
|
|
||||||
|
|
||||||
// Create a custom retryable HTTP client with specific retry configuration
|
|
||||||
customRetryClient := retryablehttp.NewClient()
|
|
||||||
customRetryClient.RetryMax = 3
|
|
||||||
customRetryClient.RetryWaitMax = 10 * time.Millisecond
|
|
||||||
|
|
||||||
// Create options with the custom retryable client
|
|
||||||
opts := defaultHTTPClientOption()
|
|
||||||
WithRetryableHTTPClint(customRetryClient)(&opts)
|
|
||||||
|
|
||||||
// Verify that the custom client is set in options
|
|
||||||
assert.NotNil(t, opts.retryableHTTPClient)
|
|
||||||
assert.Equal(t, customRetryClient, opts.retryableHTTPClient)
|
|
||||||
|
|
||||||
// Create the common client with custom retryable client
|
|
||||||
client := newCommonClient(testSystemInfo, opts)
|
|
||||||
|
|
||||||
// Make a request that will trigger a retry
|
|
||||||
req, err := http.NewRequest("GET", server.URL, nil)
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
resp, err := client.do(req)
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
// Should succeed after retry
|
|
||||||
assert.Equal(t, http.StatusOK, resp.StatusCode)
|
|
||||||
assert.Equal(t, 2, attemptCount)
|
|
||||||
|
|
||||||
// Verify that the client used is the custom one by checking newRetryableHTTPClient
|
|
||||||
retrievedRetryClient, err := client.newRetryableHTTPClient()
|
|
||||||
require.NoError(t, err)
|
|
||||||
assert.Equal(t, customRetryClient, retrievedRetryClient, "should return the custom retryable client")
|
|
||||||
})
|
|
||||||
|
|
||||||
t.Run("respects custom client's retry configuration over built-in defaults", func(t *testing.T) {
|
|
||||||
attemptCount := 0
|
|
||||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
||||||
attemptCount++
|
|
||||||
w.WriteHeader(http.StatusServiceUnavailable)
|
|
||||||
}))
|
|
||||||
defer server.Close()
|
|
||||||
|
|
||||||
// Create custom client with limited retries
|
|
||||||
customRetryClient := retryablehttp.NewClient()
|
|
||||||
customRetryClient.RetryMax = 1 // Only 1 retry (2 total attempts)
|
|
||||||
customRetryClient.RetryWaitMax = 5 * time.Millisecond
|
|
||||||
|
|
||||||
opts := defaultHTTPClientOption()
|
|
||||||
WithRetryableHTTPClint(customRetryClient)(&opts)
|
|
||||||
|
|
||||||
client := newCommonClient(testSystemInfo, opts)
|
|
||||||
|
|
||||||
req, err := http.NewRequest("GET", server.URL, nil)
|
|
||||||
require.NoError(t, err)
|
|
||||||
|
|
||||||
resp, err := client.do(req)
|
|
||||||
// When all retries are exhausted with a retryable error, the client gives up
|
|
||||||
// and an error is returned
|
|
||||||
if err != nil {
|
|
||||||
// Expected: request failed after exhausting retries
|
|
||||||
assert.Contains(t, err.Error(), "giving up after 2 attempt(s)")
|
|
||||||
} else {
|
|
||||||
// Or the final response is returned
|
|
||||||
assert.Equal(t, http.StatusServiceUnavailable, resp.StatusCode)
|
|
||||||
}
|
|
||||||
// Should have tried 1 initial + 1 retry = 2 times total
|
|
||||||
assert.Equal(t, 2, attemptCount)
|
|
||||||
})
|
|
||||||
}
|
|
||||||
@@ -74,35 +74,40 @@ func run(ctx context.Context, c Config) error {
|
|||||||
runnerGroupID = runnerGroup.ID
|
runnerGroupID = runnerGroup.ID
|
||||||
}
|
}
|
||||||
|
|
||||||
|
slog.Info("creating scale set")
|
||||||
|
|
||||||
// Create the runner scale set
|
// Create the runner scale set
|
||||||
scaleSet, err := scalesetClient.CreateRunnerScaleSet(ctx, &scaleset.RunnerScaleSet{
|
scaleSet, err := scalesetClient.CreateRunnerScaleSet(ctx, &scaleset.RunnerScaleSet{
|
||||||
Name: c.ScaleSetName,
|
Name: c.ScaleSetName,
|
||||||
RunnerGroupID: runnerGroupID,
|
RunnerGroupID: runnerGroupID,
|
||||||
Labels: c.BuildLabels(),
|
Labels: []scaleset.Label{},
|
||||||
RunnerSetting: scaleset.RunnerSetting{
|
RunnerSetting: scaleset.RunnerSetting{
|
||||||
DisableUpdate: true,
|
DisableUpdate: true,
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
slog.Error("failed to create", err)
|
||||||
return fmt.Errorf("failed to create runner scale set: %w", err)
|
return fmt.Errorf("failed to create runner scale set: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
slog.Info("created")
|
||||||
|
|
||||||
// Set the user agent for the scaleset client now that we have the scale set ID
|
// Set the user agent for the scaleset client now that we have the scale set ID
|
||||||
scalesetClient.SetSystemInfo(systemInfo(scaleSet.ID))
|
scalesetClient.SetSystemInfo(systemInfo(scaleSet.ID))
|
||||||
|
|
||||||
defer func() {
|
// defer func() {
|
||||||
logger.Info(
|
// logger.Info(
|
||||||
"Deleting runner scale set",
|
// "Deleting runner scale set",
|
||||||
slog.Int("scaleSetID", scaleSet.ID),
|
// slog.Int("scaleSetID", scaleSet.ID),
|
||||||
)
|
// )
|
||||||
if err := scalesetClient.DeleteRunnerScaleSet(context.WithoutCancel(ctx), scaleSet.ID); err != nil {
|
// if err := scalesetClient.DeleteRunnerScaleSet(context.WithoutCancel(ctx), scaleSet.ID); err != nil {
|
||||||
slog.Error(
|
// slog.Error(
|
||||||
"Failed to delete runner scale set",
|
// "Failed to delete runner scale set",
|
||||||
slog.Int("scaleSetID", scaleSet.ID),
|
// slog.Int("scaleSetID", scaleSet.ID),
|
||||||
slog.String("error", err.Error()),
|
// slog.String("error", err.Error()),
|
||||||
)
|
// )
|
||||||
}
|
// }
|
||||||
}()
|
// }()
|
||||||
|
|
||||||
dockerClient, err := dockerclient.NewClientWithOpts(dockerclient.FromEnv, dockerclient.WithAPIVersionNegotiation())
|
dockerClient, err := dockerclient.NewClientWithOpts(dockerclient.FromEnv, dockerclient.WithAPIVersionNegotiation())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
+10
-65
@@ -52,57 +52,22 @@ type Client interface {
|
|||||||
Session() scaleset.RunnerScaleSetSession
|
Session() scaleset.RunnerScaleSetSession
|
||||||
}
|
}
|
||||||
|
|
||||||
// MetricsRecorder defines the hook methods that will be called by the listener at
|
type Option func(*Listener)
|
||||||
// various points in the message handling process. This allows for custom
|
|
||||||
// metrics to be collected without coupling the listener to a specific metrics
|
|
||||||
// implementation. The methods in this interface will be called with relevant
|
|
||||||
// information about the message handling process, such as the number of jobs
|
|
||||||
// started/completed, the desired runner count, and any errors that occur.
|
|
||||||
// Implementers can use this information to track the performance and behavior
|
|
||||||
// of the listener and the scaleset service.
|
|
||||||
type MetricsRecorder interface {
|
|
||||||
RecordStatistics(statistics *scaleset.RunnerScaleSetStatistic)
|
|
||||||
RecordJobStarted(msg *scaleset.JobStarted)
|
|
||||||
RecordJobCompleted(msg *scaleset.JobCompleted)
|
|
||||||
RecordDesiredRunners(count int)
|
|
||||||
}
|
|
||||||
|
|
||||||
type discardMetricsRecorder struct{}
|
|
||||||
|
|
||||||
func (d *discardMetricsRecorder) RecordStatistics(statistics *scaleset.RunnerScaleSetStatistic) {}
|
|
||||||
func (d *discardMetricsRecorder) RecordJobStarted(msg *scaleset.JobStarted) {}
|
|
||||||
func (d *discardMetricsRecorder) RecordJobCompleted(msg *scaleset.JobCompleted) {}
|
|
||||||
func (d *discardMetricsRecorder) RecordDesiredRunners(count int) {}
|
|
||||||
|
|
||||||
// Listener listens for messages from the scaleset service and handles them. It automatically handles session
|
// Listener listens for messages from the scaleset service and handles them. It automatically handles session
|
||||||
// creation/deletion/refreshing and message polling and acking.
|
// creation/deletion/refreshing and message polling and acking.
|
||||||
type Listener struct {
|
type Listener struct {
|
||||||
// The main client responsible for communicating with the scaleset service
|
// The main client responsible for communicating with the scaleset service
|
||||||
client Client
|
client Client
|
||||||
metricsRecorder MetricsRecorder
|
|
||||||
|
|
||||||
// Configuration for the listener
|
// Configuration for the listener
|
||||||
scaleSetID int
|
scaleSetID int
|
||||||
maxRunners atomic.Uint32
|
maxRunners atomic.Uint32
|
||||||
latestStatistics *scaleset.RunnerScaleSetStatistic
|
|
||||||
|
|
||||||
// configuration for the listener
|
// configuration for the listener
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
}
|
}
|
||||||
|
|
||||||
type Option func(*Listener)
|
|
||||||
|
|
||||||
// WithMetricsRecorder sets the MetricsRecorder for the listener. If not set, a no-op recorder will be used.
|
|
||||||
// If the nil is passed, the MetricsRecorder will not be updated and the existing one will be used (which is a no-op by default).
|
|
||||||
func WithMetricsRecorder(recorder MetricsRecorder) Option {
|
|
||||||
return func(l *Listener) {
|
|
||||||
if recorder == nil {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
l.metricsRecorder = recorder
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// SetMaxRunners sets the capacity of the scaleset. It is concurrently
|
// SetMaxRunners sets the capacity of the scaleset. It is concurrently
|
||||||
// safe to update the max runners during listener.Run.
|
// safe to update the max runners during listener.Run.
|
||||||
func (l *Listener) SetMaxRunners(count int) {
|
func (l *Listener) SetMaxRunners(count int) {
|
||||||
@@ -120,17 +85,12 @@ func New(client Client, config Config, options ...Option) (*Listener, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
listener := &Listener{
|
listener := &Listener{
|
||||||
client: client,
|
client: client,
|
||||||
metricsRecorder: &discardMetricsRecorder{},
|
scaleSetID: config.ScaleSetID,
|
||||||
scaleSetID: config.ScaleSetID,
|
logger: config.Logger,
|
||||||
logger: config.Logger,
|
|
||||||
}
|
}
|
||||||
listener.SetMaxRunners(config.MaxRunners)
|
listener.SetMaxRunners(config.MaxRunners)
|
||||||
|
|
||||||
for _, option := range options {
|
|
||||||
option(listener)
|
|
||||||
}
|
|
||||||
|
|
||||||
return listener, nil
|
return listener, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -154,17 +114,13 @@ func (l *Listener) Run(ctx context.Context, scaler Scaler) error {
|
|||||||
return fmt.Errorf("session statistics is nil")
|
return fmt.Errorf("session statistics is nil")
|
||||||
}
|
}
|
||||||
|
|
||||||
l.handleStatistics(ctx, initialSession.Statistics)
|
|
||||||
|
|
||||||
l.logger.Info(
|
l.logger.Info(
|
||||||
"Handling initial session statistics",
|
"Handling initial session statistics",
|
||||||
slog.Int("totalAssignedJobs", initialSession.Statistics.TotalAssignedJobs),
|
slog.Int("totalAssignedJobs", initialSession.Statistics.TotalAssignedJobs),
|
||||||
)
|
)
|
||||||
desiredCount, err := scaler.HandleDesiredRunnerCount(ctx, initialSession.Statistics.TotalAssignedJobs)
|
if _, err := scaler.HandleDesiredRunnerCount(ctx, initialSession.Statistics.TotalAssignedJobs); err != nil {
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("handling initial message failed: %w", err)
|
return fmt.Errorf("handling initial message failed: %w", err)
|
||||||
}
|
}
|
||||||
l.metricsRecorder.RecordDesiredRunners(desiredCount)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var lastMessageID int
|
var lastMessageID int
|
||||||
@@ -186,7 +142,7 @@ func (l *Listener) Run(ctx context.Context, scaler Scaler) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if msg == nil {
|
if msg == nil {
|
||||||
_, err := scaler.HandleDesiredRunnerCount(ctx, l.latestStatistics.TotalAssignedJobs)
|
_, err := scaler.HandleDesiredRunnerCount(ctx, 0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("handling nil message failed: %w", err)
|
return fmt.Errorf("handling nil message failed: %w", err)
|
||||||
}
|
}
|
||||||
@@ -204,35 +160,24 @@ func (l *Listener) Run(ctx context.Context, scaler Scaler) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (l *Listener) handleMessage(ctx context.Context, handler Scaler, msg *scaleset.RunnerScaleSetMessage) error {
|
func (l *Listener) handleMessage(ctx context.Context, handler Scaler, msg *scaleset.RunnerScaleSetMessage) error {
|
||||||
l.handleStatistics(ctx, msg.Statistics)
|
|
||||||
|
|
||||||
if err := l.client.DeleteMessage(ctx, msg.MessageID); err != nil {
|
if err := l.client.DeleteMessage(ctx, msg.MessageID); err != nil {
|
||||||
return fmt.Errorf("failed to delete message: %w", err)
|
return fmt.Errorf("failed to delete message: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
for _, jobStarted := range msg.JobStartedMessages {
|
for _, jobStarted := range msg.JobStartedMessages {
|
||||||
l.metricsRecorder.RecordJobStarted(jobStarted)
|
|
||||||
if err := handler.HandleJobStarted(ctx, jobStarted); err != nil {
|
if err := handler.HandleJobStarted(ctx, jobStarted); err != nil {
|
||||||
return fmt.Errorf("failed to handle job started: %w", err)
|
return fmt.Errorf("failed to handle job started: %w", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
for _, jobCompleted := range msg.JobCompletedMessages {
|
for _, jobCompleted := range msg.JobCompletedMessages {
|
||||||
l.metricsRecorder.RecordJobCompleted(jobCompleted)
|
|
||||||
if err := handler.HandleJobCompleted(ctx, jobCompleted); err != nil {
|
if err := handler.HandleJobCompleted(ctx, jobCompleted); err != nil {
|
||||||
return fmt.Errorf("failed to handle job completed: %w", err)
|
return fmt.Errorf("failed to handle job completed: %w", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
desiredCount, err := handler.HandleDesiredRunnerCount(ctx, msg.Statistics.TotalAssignedJobs)
|
if _, err := handler.HandleDesiredRunnerCount(ctx, msg.Statistics.TotalAssignedJobs); err != nil {
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("failed to handle desired runner count: %w", err)
|
return fmt.Errorf("failed to handle desired runner count: %w", err)
|
||||||
}
|
}
|
||||||
l.metricsRecorder.RecordDesiredRunners(desiredCount)
|
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (l *Listener) handleStatistics(ctx context.Context, msg *scaleset.RunnerScaleSetStatistic) {
|
|
||||||
l.latestStatistics = msg
|
|
||||||
l.metricsRecorder.RecordStatistics(msg)
|
|
||||||
}
|
|
||||||
+51
-48
@@ -75,37 +75,39 @@ func TestListener_Run(t *testing.T) {
|
|||||||
client := NewMockClient(t)
|
client := NewMockClient(t)
|
||||||
|
|
||||||
uuid := uuid.New()
|
uuid := uuid.New()
|
||||||
initialStatistics := &scaleset.RunnerScaleSetStatistic{
|
|
||||||
TotalAssignedJobs: 2,
|
|
||||||
}
|
|
||||||
session := scaleset.RunnerScaleSetSession{
|
session := scaleset.RunnerScaleSetSession{
|
||||||
SessionID: uuid,
|
SessionID: uuid,
|
||||||
OwnerName: "example",
|
OwnerName: "example",
|
||||||
RunnerScaleSet: &scaleset.RunnerScaleSet{},
|
RunnerScaleSet: &scaleset.RunnerScaleSet{},
|
||||||
MessageQueueURL: "https://example.com",
|
MessageQueueURL: "https://example.com",
|
||||||
MessageQueueAccessToken: "1234567890",
|
MessageQueueAccessToken: "1234567890",
|
||||||
Statistics: initialStatistics,
|
Statistics: &scaleset.RunnerScaleSetStatistic{},
|
||||||
}
|
}
|
||||||
|
|
||||||
client.On("Session").Return(session).Once()
|
client.On("Session").Return(session).Once()
|
||||||
|
|
||||||
metricsRecorder := NewMockMetricsRecorder(t)
|
l, err := New(client, config)
|
||||||
metricsRecorder.On("RecordStatistics", initialStatistics).Once()
|
|
||||||
metricsRecorder.On("RecordDesiredRunners", initialStatistics.TotalAssignedJobs).
|
|
||||||
Return(initialStatistics.TotalAssignedJobs, nil).
|
|
||||||
Run(func(mock.Arguments) { cancel() }).
|
|
||||||
Once()
|
|
||||||
|
|
||||||
l, err := New(client, config, WithMetricsRecorder(metricsRecorder))
|
|
||||||
require.Nil(t, err)
|
require.Nil(t, err)
|
||||||
|
|
||||||
|
var called bool
|
||||||
handler := NewMockScaler(t)
|
handler := NewMockScaler(t)
|
||||||
handler.On("HandleDesiredRunnerCount", mock.Anything, mock.Anything).
|
handler.On(
|
||||||
Return(initialStatistics.TotalAssignedJobs, nil).
|
"HandleDesiredRunnerCount",
|
||||||
|
mock.Anything,
|
||||||
|
mock.Anything,
|
||||||
|
).
|
||||||
|
Return(0, nil).
|
||||||
|
Run(
|
||||||
|
func(mock.Arguments) {
|
||||||
|
called = true
|
||||||
|
cancel()
|
||||||
|
},
|
||||||
|
).
|
||||||
Once()
|
Once()
|
||||||
|
|
||||||
err = l.Run(ctx, handler)
|
err = l.Run(ctx, handler)
|
||||||
assert.ErrorIs(t, err, context.Canceled)
|
assert.ErrorIs(t, err, context.Canceled)
|
||||||
|
assert.True(t, called)
|
||||||
})
|
})
|
||||||
|
|
||||||
t.Run("cancel context after get message", func(t *testing.T) {
|
t.Run("cancel context after get message", func(t *testing.T) {
|
||||||
@@ -118,60 +120,61 @@ func TestListener_Run(t *testing.T) {
|
|||||||
MaxRunners: 10,
|
MaxRunners: 10,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
client := NewMockClient(t)
|
||||||
uuid := uuid.New()
|
uuid := uuid.New()
|
||||||
initialStatistics := &scaleset.RunnerScaleSetStatistic{
|
|
||||||
TotalAssignedJobs: 2,
|
|
||||||
}
|
|
||||||
|
|
||||||
session := scaleset.RunnerScaleSetSession{
|
session := scaleset.RunnerScaleSetSession{
|
||||||
SessionID: uuid,
|
SessionID: uuid,
|
||||||
OwnerName: "example",
|
OwnerName: "example",
|
||||||
RunnerScaleSet: &scaleset.RunnerScaleSet{},
|
RunnerScaleSet: &scaleset.RunnerScaleSet{},
|
||||||
MessageQueueURL: "https://example.com",
|
MessageQueueURL: "https://example.com",
|
||||||
MessageQueueAccessToken: "1234567890",
|
MessageQueueAccessToken: "1234567890",
|
||||||
Statistics: initialStatistics,
|
Statistics: &scaleset.RunnerScaleSetStatistic{},
|
||||||
}
|
}
|
||||||
|
|
||||||
msg := &scaleset.RunnerScaleSetMessage{
|
msg := &scaleset.RunnerScaleSetMessage{
|
||||||
MessageID: 1,
|
MessageID: 1,
|
||||||
Statistics: &scaleset.RunnerScaleSetStatistic{
|
Statistics: &scaleset.RunnerScaleSetStatistic{},
|
||||||
TotalAssignedJobs: 3,
|
|
||||||
},
|
|
||||||
}
|
}
|
||||||
|
|
||||||
metricsRecorder := NewMockMetricsRecorder(t)
|
|
||||||
client := NewMockClient(t)
|
|
||||||
handler := NewMockScaler(t)
|
|
||||||
|
|
||||||
client.On("Session").Return(session).Once()
|
client.On("Session").Return(session).Once()
|
||||||
metricsRecorder.On("RecordStatistics", initialStatistics).Once()
|
client.On(
|
||||||
metricsRecorder.On("RecordDesiredRunners", initialStatistics.TotalAssignedJobs).
|
"GetMessage",
|
||||||
Return(initialStatistics.TotalAssignedJobs, nil).
|
ctx,
|
||||||
Once()
|
mock.Anything,
|
||||||
handler.On("HandleDesiredRunnerCount", mock.Anything, initialStatistics.TotalAssignedJobs).
|
10,
|
||||||
Return(initialStatistics.TotalAssignedJobs, nil).
|
).
|
||||||
Once()
|
|
||||||
|
|
||||||
client.On("GetMessage", ctx, mock.Anything, 10).
|
|
||||||
Return(msg, nil).
|
Return(msg, nil).
|
||||||
Run(func(mock.Arguments) { cancel() }).
|
Run(
|
||||||
|
func(mock.Arguments) {
|
||||||
|
cancel()
|
||||||
|
},
|
||||||
|
).
|
||||||
Once()
|
Once()
|
||||||
|
|
||||||
metricsRecorder.On("RecordStatistics", msg.Statistics).Once()
|
|
||||||
// Ensure delete message is called without cancel
|
// Ensure delete message is called without cancel
|
||||||
client.On("DeleteMessage", context.WithoutCancel(ctx), mock.Anything).
|
client.On(
|
||||||
Return(nil).
|
"DeleteMessage",
|
||||||
|
context.WithoutCancel(ctx),
|
||||||
|
mock.Anything,
|
||||||
|
).Return(nil).Once()
|
||||||
|
|
||||||
|
handler := NewMockScaler(t)
|
||||||
|
handler.On(
|
||||||
|
"HandleDesiredRunnerCount",
|
||||||
|
mock.Anything,
|
||||||
|
0,
|
||||||
|
).
|
||||||
|
Return(0, nil).
|
||||||
Once()
|
Once()
|
||||||
|
|
||||||
metricsRecorder.On("RecordDesiredRunners", msg.Statistics.TotalAssignedJobs).
|
handler.On(
|
||||||
Return(msg.Statistics.TotalAssignedJobs, nil).
|
"HandleDesiredRunnerCount",
|
||||||
|
mock.Anything,
|
||||||
|
mock.Anything,
|
||||||
|
).
|
||||||
|
Return(0, nil).
|
||||||
Once()
|
Once()
|
||||||
|
|
||||||
handler.On("HandleDesiredRunnerCount", mock.Anything, msg.Statistics.TotalAssignedJobs).
|
l, err := New(client, config)
|
||||||
Return(msg.Statistics.TotalAssignedJobs, nil).
|
|
||||||
Once()
|
|
||||||
|
|
||||||
l, err := New(client, config, WithMetricsRecorder(metricsRecorder))
|
|
||||||
require.Nil(t, err)
|
require.Nil(t, err)
|
||||||
|
|
||||||
err = l.Run(ctx, handler)
|
err = l.Run(ctx, handler)
|
||||||
|
|||||||
@@ -213,193 +213,6 @@ func (_c *MockClient_Session_Call) RunAndReturn(run func() scaleset.RunnerScaleS
|
|||||||
return _c
|
return _c
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewMockMetricsRecorder creates a new instance of MockMetricsRecorder. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations.
|
|
||||||
// The first argument is typically a *testing.T value.
|
|
||||||
func NewMockMetricsRecorder(t interface {
|
|
||||||
mock.TestingT
|
|
||||||
Cleanup(func())
|
|
||||||
}) *MockMetricsRecorder {
|
|
||||||
mock := &MockMetricsRecorder{}
|
|
||||||
mock.Mock.Test(t)
|
|
||||||
|
|
||||||
t.Cleanup(func() { mock.AssertExpectations(t) })
|
|
||||||
|
|
||||||
return mock
|
|
||||||
}
|
|
||||||
|
|
||||||
// MockMetricsRecorder is an autogenerated mock type for the MetricsRecorder type
|
|
||||||
type MockMetricsRecorder struct {
|
|
||||||
mock.Mock
|
|
||||||
}
|
|
||||||
|
|
||||||
type MockMetricsRecorder_Expecter struct {
|
|
||||||
mock *mock.Mock
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_m *MockMetricsRecorder) EXPECT() *MockMetricsRecorder_Expecter {
|
|
||||||
return &MockMetricsRecorder_Expecter{mock: &_m.Mock}
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordDesiredRunners provides a mock function for the type MockMetricsRecorder
|
|
||||||
func (_mock *MockMetricsRecorder) RecordDesiredRunners(count int) {
|
|
||||||
_mock.Called(count)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// MockMetricsRecorder_RecordDesiredRunners_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'RecordDesiredRunners'
|
|
||||||
type MockMetricsRecorder_RecordDesiredRunners_Call struct {
|
|
||||||
*mock.Call
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordDesiredRunners is a helper method to define mock.On call
|
|
||||||
// - count int
|
|
||||||
func (_e *MockMetricsRecorder_Expecter) RecordDesiredRunners(count interface{}) *MockMetricsRecorder_RecordDesiredRunners_Call {
|
|
||||||
return &MockMetricsRecorder_RecordDesiredRunners_Call{Call: _e.mock.On("RecordDesiredRunners", count)}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordDesiredRunners_Call) Run(run func(count int)) *MockMetricsRecorder_RecordDesiredRunners_Call {
|
|
||||||
_c.Call.Run(func(args mock.Arguments) {
|
|
||||||
var arg0 int
|
|
||||||
if args[0] != nil {
|
|
||||||
arg0 = args[0].(int)
|
|
||||||
}
|
|
||||||
run(
|
|
||||||
arg0,
|
|
||||||
)
|
|
||||||
})
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordDesiredRunners_Call) Return() *MockMetricsRecorder_RecordDesiredRunners_Call {
|
|
||||||
_c.Call.Return()
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordDesiredRunners_Call) RunAndReturn(run func(count int)) *MockMetricsRecorder_RecordDesiredRunners_Call {
|
|
||||||
_c.Run(run)
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordJobCompleted provides a mock function for the type MockMetricsRecorder
|
|
||||||
func (_mock *MockMetricsRecorder) RecordJobCompleted(msg *scaleset.JobCompleted) {
|
|
||||||
_mock.Called(msg)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// MockMetricsRecorder_RecordJobCompleted_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'RecordJobCompleted'
|
|
||||||
type MockMetricsRecorder_RecordJobCompleted_Call struct {
|
|
||||||
*mock.Call
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordJobCompleted is a helper method to define mock.On call
|
|
||||||
// - msg *scaleset.JobCompleted
|
|
||||||
func (_e *MockMetricsRecorder_Expecter) RecordJobCompleted(msg interface{}) *MockMetricsRecorder_RecordJobCompleted_Call {
|
|
||||||
return &MockMetricsRecorder_RecordJobCompleted_Call{Call: _e.mock.On("RecordJobCompleted", msg)}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordJobCompleted_Call) Run(run func(msg *scaleset.JobCompleted)) *MockMetricsRecorder_RecordJobCompleted_Call {
|
|
||||||
_c.Call.Run(func(args mock.Arguments) {
|
|
||||||
var arg0 *scaleset.JobCompleted
|
|
||||||
if args[0] != nil {
|
|
||||||
arg0 = args[0].(*scaleset.JobCompleted)
|
|
||||||
}
|
|
||||||
run(
|
|
||||||
arg0,
|
|
||||||
)
|
|
||||||
})
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordJobCompleted_Call) Return() *MockMetricsRecorder_RecordJobCompleted_Call {
|
|
||||||
_c.Call.Return()
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordJobCompleted_Call) RunAndReturn(run func(msg *scaleset.JobCompleted)) *MockMetricsRecorder_RecordJobCompleted_Call {
|
|
||||||
_c.Run(run)
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordJobStarted provides a mock function for the type MockMetricsRecorder
|
|
||||||
func (_mock *MockMetricsRecorder) RecordJobStarted(msg *scaleset.JobStarted) {
|
|
||||||
_mock.Called(msg)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// MockMetricsRecorder_RecordJobStarted_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'RecordJobStarted'
|
|
||||||
type MockMetricsRecorder_RecordJobStarted_Call struct {
|
|
||||||
*mock.Call
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordJobStarted is a helper method to define mock.On call
|
|
||||||
// - msg *scaleset.JobStarted
|
|
||||||
func (_e *MockMetricsRecorder_Expecter) RecordJobStarted(msg interface{}) *MockMetricsRecorder_RecordJobStarted_Call {
|
|
||||||
return &MockMetricsRecorder_RecordJobStarted_Call{Call: _e.mock.On("RecordJobStarted", msg)}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordJobStarted_Call) Run(run func(msg *scaleset.JobStarted)) *MockMetricsRecorder_RecordJobStarted_Call {
|
|
||||||
_c.Call.Run(func(args mock.Arguments) {
|
|
||||||
var arg0 *scaleset.JobStarted
|
|
||||||
if args[0] != nil {
|
|
||||||
arg0 = args[0].(*scaleset.JobStarted)
|
|
||||||
}
|
|
||||||
run(
|
|
||||||
arg0,
|
|
||||||
)
|
|
||||||
})
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordJobStarted_Call) Return() *MockMetricsRecorder_RecordJobStarted_Call {
|
|
||||||
_c.Call.Return()
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordJobStarted_Call) RunAndReturn(run func(msg *scaleset.JobStarted)) *MockMetricsRecorder_RecordJobStarted_Call {
|
|
||||||
_c.Run(run)
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordStatistics provides a mock function for the type MockMetricsRecorder
|
|
||||||
func (_mock *MockMetricsRecorder) RecordStatistics(statistics *scaleset.RunnerScaleSetStatistic) {
|
|
||||||
_mock.Called(statistics)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// MockMetricsRecorder_RecordStatistics_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'RecordStatistics'
|
|
||||||
type MockMetricsRecorder_RecordStatistics_Call struct {
|
|
||||||
*mock.Call
|
|
||||||
}
|
|
||||||
|
|
||||||
// RecordStatistics is a helper method to define mock.On call
|
|
||||||
// - statistics *scaleset.RunnerScaleSetStatistic
|
|
||||||
func (_e *MockMetricsRecorder_Expecter) RecordStatistics(statistics interface{}) *MockMetricsRecorder_RecordStatistics_Call {
|
|
||||||
return &MockMetricsRecorder_RecordStatistics_Call{Call: _e.mock.On("RecordStatistics", statistics)}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordStatistics_Call) Run(run func(statistics *scaleset.RunnerScaleSetStatistic)) *MockMetricsRecorder_RecordStatistics_Call {
|
|
||||||
_c.Call.Run(func(args mock.Arguments) {
|
|
||||||
var arg0 *scaleset.RunnerScaleSetStatistic
|
|
||||||
if args[0] != nil {
|
|
||||||
arg0 = args[0].(*scaleset.RunnerScaleSetStatistic)
|
|
||||||
}
|
|
||||||
run(
|
|
||||||
arg0,
|
|
||||||
)
|
|
||||||
})
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordStatistics_Call) Return() *MockMetricsRecorder_RecordStatistics_Call {
|
|
||||||
_c.Call.Return()
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
func (_c *MockMetricsRecorder_RecordStatistics_Call) RunAndReturn(run func(statistics *scaleset.RunnerScaleSetStatistic)) *MockMetricsRecorder_RecordStatistics_Call {
|
|
||||||
_c.Run(run)
|
|
||||||
return _c
|
|
||||||
}
|
|
||||||
|
|
||||||
// NewMockScaler creates a new instance of MockScaler. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations.
|
// NewMockScaler creates a new instance of MockScaler. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations.
|
||||||
// The first argument is typically a *testing.T value.
|
// The first argument is typically a *testing.T value.
|
||||||
func NewMockScaler(t interface {
|
func NewMockScaler(t interface {
|
||||||
|
|||||||
@@ -235,9 +235,6 @@ func (c *MessageSessionClient) Session() RunnerScaleSetSession {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *MessageSessionClient) doSessionRequest(ctx context.Context, method, path string, requestData io.Reader, expectedResponseStatusCode int, responseUnmarshalTarget any) error {
|
func (c *MessageSessionClient) doSessionRequest(ctx context.Context, method, path string, requestData io.Reader, expectedResponseStatusCode int, responseUnmarshalTarget any) error {
|
||||||
c.innerClient.mu.Lock()
|
|
||||||
defer c.innerClient.mu.Unlock()
|
|
||||||
|
|
||||||
req, err := c.innerClient.newActionsServiceRequest(ctx, method, path, requestData)
|
req, err := c.innerClient.newActionsServiceRequest(ctx, method, path, requestData)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("failed to create new actions service request: %w", err)
|
return fmt.Errorf("failed to create new actions service request: %w", err)
|
||||||
|
|||||||
Reference in new issue
Block a user