mirror of
https://gitea.com/gitea/act_runner.git
synced 2026-08-13 22:41:54 +02:00
fix(cache): build job URLs on the address its runner registered (#1153)
The cache server built every URL it hands a job from its own listen address, so jobs whose runner reaches it through a reverse proxy were sent to the internal one. This covered the v1 `archiveLocation`, the v2 signed cache URLs and `ACTIONS_RESULTS_URL`. Runners now register the address their jobs reach the server at, next to the instance URL they already send. The cache-server needs no configuration of its own, and runners that reach it differently each get their own correct address. Closes https://gitea.com/gitea/runner/issues/1152 --------- Co-authored-by: silverwind <[email protected]> Reviewed-on: https://gitea.com/gitea/runner/pulls/1153 Reviewed-by: silverwind <[email protected]> Reviewed-by: bircni <[email protected]> Co-authored-by: Max P. <[email protected]>
This commit is contained in:
committed by
bircni
co-authored by
silverwind
parent
b66433e667
commit
e178c03adc
@@ -313,6 +313,8 @@ Run one dedicated `gitea-runner cache-server` that all runners point at.
|
|||||||
# external_secret_file: /path/to/secret # secret can also be passed via a file
|
# external_secret_file: /path/to/secret # secret can also be passed via a file
|
||||||
```
|
```
|
||||||
|
|
||||||
|
Jobs reach the cache server at `external_server`, so when a reverse proxy fronts the server, point `external_server` at the proxy. The cache server itself needs no extra configuration.
|
||||||
|
|
||||||
Alternatively, mount the same NFS/CIFS share on every runner and point `cache.dir` at it — simpler, but with weaker isolation between repositories.
|
Alternatively, mount the same NFS/CIFS share on every runner and point `cache.dir` at it — simpler, but with weaker isolation between repositories.
|
||||||
|
|
||||||
**S3 / MinIO** — mount object storage as a FUSE filesystem (e.g. [s3fs](https://github.com/s3fs-fuse/s3fs-fuse) or [goofys](https://github.com/kahing/goofys)) and set `cache.dir` to the mount point.
|
**S3 / MinIO** — mount object storage as a FUSE filesystem (e.g. [s3fs](https://github.com/s3fs-fuse/s3fs-fuse) or [goofys](https://github.com/kahing/goofys)) and set `cache.dir` to the mount point.
|
||||||
|
|||||||
@@ -59,6 +59,9 @@ type JobCredential struct {
|
|||||||
// remote runner registers with.
|
// remote runner registers with.
|
||||||
Results string `json:"results"`
|
Results string `json:"results"`
|
||||||
InsecureTLS bool `json:"insecure_tls"`
|
InsecureTLS bool `json:"insecure_tls"`
|
||||||
|
|
||||||
|
// PublicURL is this server as a reverse proxy makes the job reach it, not the listen address.
|
||||||
|
PublicURL string `json:"public_url"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// credEntry holds a registered job's credential along with an active
|
// credEntry holds a registered job's credential along with an active
|
||||||
@@ -212,6 +215,13 @@ func (h *Handler) ExternalURL() string {
|
|||||||
return fmt.Sprintf("http://%s:%d", h.outboundIP, h.port)
|
return fmt.Sprintf("http://%s:%d", h.outboundIP, h.port)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (h *Handler) baseURL(cred JobCredential) string {
|
||||||
|
if base := strings.TrimRight(cred.PublicURL, "/"); base != "" {
|
||||||
|
return base
|
||||||
|
}
|
||||||
|
return h.ExternalURL()
|
||||||
|
}
|
||||||
|
|
||||||
// RegisterJob makes token a valid bearer credential for cache requests from
|
// RegisterJob makes token a valid bearer credential for cache requests from
|
||||||
// the given repository and returns a function that removes it. The runner
|
// the given repository and returns a function that removes it. The runner
|
||||||
// calls this at job start and defers the returned func so that the credential
|
// calls this at job start and defers the returned func so that the credential
|
||||||
@@ -359,7 +369,7 @@ func (h *Handler) find(w http.ResponseWriter, r *http.Request, _ httprouter.Para
|
|||||||
}
|
}
|
||||||
h.responseJSON(w, r, 200, map[string]any{
|
h.responseJSON(w, r, 200, map[string]any{
|
||||||
"result": "hit",
|
"result": "hit",
|
||||||
"archiveLocation": h.signedArtifactURL(cache.ID, time.Now().Add(artifactURLTTL)),
|
"archiveLocation": h.signedArtifactURL(cred, cache.ID, time.Now().Add(artifactURLTTL)),
|
||||||
"cacheKey": cache.Key,
|
"cacheKey": cache.Key,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -641,7 +651,7 @@ func (h *Handler) ResultsURL(cred JobCredential) string {
|
|||||||
if h == nil || cred.Results == "" {
|
if h == nil || cred.Results == "" {
|
||||||
return ""
|
return ""
|
||||||
}
|
}
|
||||||
return h.ExternalURL()
|
return h.baseURL(cred)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (h *Handler) internalRegister(w http.ResponseWriter, r *http.Request, _ httprouter.Params) {
|
func (h *Handler) internalRegister(w http.ResponseWriter, r *http.Request, _ httprouter.Params) {
|
||||||
@@ -700,16 +710,16 @@ func (h *Handler) computeSignature(purpose string, cacheID, exp int64) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// signedURL builds a URL under path that signedAuth accepts for the same purpose.
|
// signedURL builds a URL under path that signedAuth accepts for the same purpose.
|
||||||
func (h *Handler) signedURL(path, purpose string, cacheID uint64, exp time.Time) string {
|
func (h *Handler) signedURL(cred JobCredential, path, purpose string, cacheID uint64, exp time.Time) string {
|
||||||
expUnix := exp.Unix()
|
expUnix := exp.Unix()
|
||||||
q := url.Values{}
|
q := url.Values{}
|
||||||
q.Set("exp", strconv.FormatInt(expUnix, 10))
|
q.Set("exp", strconv.FormatInt(expUnix, 10))
|
||||||
q.Set("sig", h.computeSignature(purpose, int64(cacheID), expUnix))
|
q.Set("sig", h.computeSignature(purpose, int64(cacheID), expUnix))
|
||||||
return fmt.Sprintf("%s%s/%d?%s", h.ExternalURL(), path, cacheID, q.Encode())
|
return fmt.Sprintf("%s%s/%d?%s", h.baseURL(cred), path, cacheID, q.Encode())
|
||||||
}
|
}
|
||||||
|
|
||||||
func (h *Handler) signedArtifactURL(cacheID uint64, exp time.Time) string {
|
func (h *Handler) signedArtifactURL(cred JobCredential, cacheID uint64, exp time.Time) string {
|
||||||
return h.signedURL(apiPath+"/artifacts", "", cacheID, exp)
|
return h.signedURL(cred, apiPath+"/artifacts", "", cacheID, exp)
|
||||||
}
|
}
|
||||||
|
|
||||||
// if not found, return (nil, nil) instead of an error.
|
// if not found, return (nil, nil) instead of an error.
|
||||||
|
|||||||
@@ -45,7 +45,7 @@ var testClient = &http.Client{Transport: &bearerTransport{token: testToken}}
|
|||||||
// tests use it to reach the get handler directly without going through a
|
// tests use it to reach the get handler directly without going through a
|
||||||
// find/cache-hit round trip.
|
// find/cache-hit round trip.
|
||||||
func signArtifactURL(h *Handler, id int64) string {
|
func signArtifactURL(h *Handler, id int64) string {
|
||||||
return h.signedArtifactURL(uint64(id), time.Now().Add(artifactURLTTL))
|
return h.signedArtifactURL(JobCredential{}, uint64(id), time.Now().Add(artifactURLTTL))
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestHandler(t *testing.T) {
|
func TestHandler(t *testing.T) {
|
||||||
@@ -998,7 +998,7 @@ func TestHandler_ArtifactSignature(t *testing.T) {
|
|||||||
})
|
})
|
||||||
|
|
||||||
t.Run("tampered signature", func(t *testing.T) {
|
t.Run("tampered signature", func(t *testing.T) {
|
||||||
good := handler.signedArtifactURL(1, time.Now().Add(artifactURLTTL))
|
good := signArtifactURL(handler, 1)
|
||||||
bad := good[:len(good)-4] + "dead"
|
bad := good[:len(good)-4] + "dead"
|
||||||
resp, err := testClient.Get(bad)
|
resp, err := testClient.Get(bad)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
@@ -1007,7 +1007,7 @@ func TestHandler_ArtifactSignature(t *testing.T) {
|
|||||||
})
|
})
|
||||||
|
|
||||||
t.Run("expired signature", func(t *testing.T) {
|
t.Run("expired signature", func(t *testing.T) {
|
||||||
expired := handler.signedArtifactURL(1, time.Now().Add(-time.Second))
|
expired := handler.signedArtifactURL(JobCredential{}, 1, time.Now().Add(-time.Second))
|
||||||
resp, err := testClient.Get(expired)
|
resp, err := testClient.Get(expired)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
resp.Body.Close()
|
resp.Body.Close()
|
||||||
@@ -1019,7 +1019,7 @@ func TestHandler_ArtifactSignature(t *testing.T) {
|
|||||||
other, err := StartHandler(dir2, "", 0, "", nil)
|
other, err := StartHandler(dir2, "", 0, "", nil)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
defer other.Close()
|
defer other.Close()
|
||||||
otherURL := other.signedArtifactURL(1, time.Now().Add(artifactURLTTL))
|
otherURL := signArtifactURL(other, 1)
|
||||||
// Rewrite the host so the request still lands on our handler, but
|
// Rewrite the host so the request still lands on our handler, but
|
||||||
// the signature was computed with a different secret.
|
// the signature was computed with a different secret.
|
||||||
parts := strings.SplitN(otherURL, apiPath, 2)
|
parts := strings.SplitN(otherURL, apiPath, 2)
|
||||||
|
|||||||
@@ -97,7 +97,7 @@ func (h *Handler) v2CreateCacheEntry(w http.ResponseWriter, r *http.Request, _ h
|
|||||||
|
|
||||||
h.responseJSON(w, r, http.StatusOK, map[string]any{
|
h.responseJSON(w, r, http.StatusOK, map[string]any{
|
||||||
"ok": true,
|
"ok": true,
|
||||||
"signed_upload_url": h.signedURL(blobPath, blobUploadPurpose, cache.ID, time.Now().Add(blobUploadURLTTL)),
|
"signed_upload_url": h.signedURL(cred, blobPath, blobUploadPurpose, cache.ID, time.Now().Add(blobUploadURLTTL)),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -168,7 +168,7 @@ func (h *Handler) v2GetCacheEntryDownloadURL(w http.ResponseWriter, r *http.Requ
|
|||||||
|
|
||||||
h.responseJSON(w, r, http.StatusOK, map[string]any{
|
h.responseJSON(w, r, http.StatusOK, map[string]any{
|
||||||
"ok": true,
|
"ok": true,
|
||||||
"signed_download_url": h.signedArtifactURL(cache.ID, time.Now().Add(artifactURLTTL)),
|
"signed_download_url": h.signedArtifactURL(cred, cache.ID, time.Now().Add(artifactURLTTL)),
|
||||||
"matched_key": cache.Key,
|
"matched_key": cache.Key,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
@@ -198,6 +199,18 @@ func TestCacheServiceV2Lookups(t *testing.T) {
|
|||||||
assert.NotEmpty(t, reserved["signed_upload_url"])
|
assert.NotEmpty(t, reserved["signed_upload_url"])
|
||||||
})
|
})
|
||||||
|
|
||||||
|
t.Run("a proxied job is handed the address its runner registered", func(t *testing.T) {
|
||||||
|
const proxy = "https://cache.example.invalid"
|
||||||
|
handler.RegisterJob("proxied", JobCredential{Repo: testRepo, PublicURL: proxy + "/"})
|
||||||
|
client := &http.Client{Transport: &bearerTransport{token: "proxied"}}
|
||||||
|
|
||||||
|
created := v2Call(t, handler, client, "CreateCacheEntry", map[string]any{"key": "proxied-key", "version": "v1"})
|
||||||
|
assert.True(t, strings.HasPrefix(created["signed_upload_url"].(string), proxy+blobPath+"/"))
|
||||||
|
|
||||||
|
got := v2Call(t, handler, client, "GetCacheEntryDownloadURL", map[string]any{"key": "deps-abc", "version": "v1"})
|
||||||
|
assert.True(t, strings.HasPrefix(got["signed_download_url"].(string), proxy+apiPath+"/artifacts/"))
|
||||||
|
})
|
||||||
|
|
||||||
t.Run("finalizing without a reservation is not ok", func(t *testing.T) {
|
t.Run("finalizing without a reservation is not ok", func(t *testing.T) {
|
||||||
got := v2Call(t, handler, testClient, "FinalizeCacheEntryUpload", map[string]any{
|
got := v2Call(t, handler, testClient, "FinalizeCacheEntryUpload", map[string]any{
|
||||||
"key": "never-reserved", "version": "v1", "size_bytes": 1,
|
"key": "never-reserved", "version": "v1", "size_bytes": 1,
|
||||||
|
|||||||
@@ -605,14 +605,15 @@ func (r *Runner) registerExternalCacheJob(token string, cred artifactcache.JobCr
|
|||||||
resultsURL := ""
|
resultsURL := ""
|
||||||
if body, err := postInternalCache(base+"/_internal/register", r.cfg.Cache.ExternalSecret, map[string]any{
|
if body, err := postInternalCache(base+"/_internal/register", r.cfg.Cache.ExternalSecret, map[string]any{
|
||||||
"token": token, "repo": cred.Repo, "results": cred.Results, "insecure_tls": cred.InsecureTLS,
|
"token": token, "repo": cred.Repo, "results": cred.Results, "insecure_tls": cred.InsecureTLS,
|
||||||
|
"public_url": base,
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
log.Warnf("cache external_server register failed (%s): %v", base, err)
|
log.Warnf("cache external_server register failed (%s): %v", base, err)
|
||||||
if reporter != nil {
|
if reporter != nil {
|
||||||
reporter.Logf("::warning::%s", runner.EscapeCommandData(fmt.Sprintf(
|
reporter.Logf("::warning::%s", runner.EscapeCommandData(fmt.Sprintf(
|
||||||
"cache external_server register failed (%s): %v — cache requests from this job will be unauthenticated and likely return 401", base, err)))
|
"cache external_server register failed (%s): %v — cache requests from this job will be unauthenticated and likely return 401", base, err)))
|
||||||
}
|
}
|
||||||
} else {
|
} else if forwarded, _ := body["results_url"].(string); forwarded != "" {
|
||||||
resultsURL, _ = body["results_url"].(string) // absent from a server too old to forward
|
resultsURL = base // the answer only says it forwards, its own address need not be the job's
|
||||||
}
|
}
|
||||||
return func() {
|
return func() {
|
||||||
if _, err := postInternalCache(base+"/_internal/revoke", r.cfg.Cache.ExternalSecret,
|
if _, err := postInternalCache(base+"/_internal/revoke", r.cfg.Cache.ExternalSecret,
|
||||||
|
|||||||
@@ -157,14 +157,15 @@ func decodeJSON(resp *http.Response, v any) error {
|
|||||||
|
|
||||||
// End-to-end against a remote cache-server: token unknown → 401, register →
|
// End-to-end against a remote cache-server: token unknown → 401, register →
|
||||||
// reserve/upload/commit/find/download all OK, revoke → 401 again. Registering also names the
|
// reserve/upload/commit/find/download all OK, revoke → 401 again. Registering also names the
|
||||||
// instance, and the server answering with its own address is what makes a shared cache server the
|
// instance and the address its jobs reach the server at, which makes a shared cache server the
|
||||||
// whole results service, as the built-in one is.
|
// whole results service, as the built-in one is.
|
||||||
func TestRunner_ExternalCacheServer_RegisterRevoke(t *testing.T) {
|
func TestRunner_ExternalCacheServer_RegisterRevoke(t *testing.T) {
|
||||||
dir := filepath.Join(t.TempDir(), "remote-cache")
|
dir := filepath.Join(t.TempDir(), "remote-cache")
|
||||||
const secret = "shared-secret-for-tests"
|
const secret = "shared-secret-for-tests"
|
||||||
remote, err := artifactcache.StartHandler(dir, "127.0.0.1", 0, secret, nil)
|
remote, err := artifactcache.StartHandler(dir, "127.0.0.2", 0, secret, nil) // advertised, never dialled
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
defer remote.Close()
|
defer remote.Close()
|
||||||
|
external := strings.Replace(remote.ExternalURL(), "127.0.0.2", "127.0.0.1", 1)
|
||||||
gitea := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
gitea := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||||
_, _ = io.WriteString(w, `{"ok":true}`)
|
_, _ = io.WriteString(w, `{"ok":true}`)
|
||||||
}))
|
}))
|
||||||
@@ -172,7 +173,7 @@ func TestRunner_ExternalCacheServer_RegisterRevoke(t *testing.T) {
|
|||||||
|
|
||||||
r := &Runner{
|
r := &Runner{
|
||||||
cfg: &config.Config{Cache: config.Cache{
|
cfg: &config.Config{Cache: config.Cache{
|
||||||
ExternalServer: remote.ExternalURL(),
|
ExternalServer: external,
|
||||||
ExternalSecret: secret,
|
ExternalSecret: secret,
|
||||||
}},
|
}},
|
||||||
envs: map[string]string{"ACTIONS_RESULTS_URL": gitea.URL},
|
envs: map[string]string{"ACTIONS_RESULTS_URL": gitea.URL},
|
||||||
@@ -180,7 +181,7 @@ func TestRunner_ExternalCacheServer_RegisterRevoke(t *testing.T) {
|
|||||||
|
|
||||||
token := "external-task-token"
|
token := "external-task-token"
|
||||||
repo := "owner/repoX"
|
repo := "owner/repoX"
|
||||||
base := remote.ExternalURL() + "/_apis/artifactcache"
|
base := external + "/_apis/artifactcache"
|
||||||
probe := func() int {
|
probe := func() int {
|
||||||
req, _ := http.NewRequest(http.MethodGet, base+"/cache?keys=k&version=v", nil)
|
req, _ := http.NewRequest(http.MethodGet, base+"/cache?keys=k&version=v", nil)
|
||||||
req.Header.Set("Authorization", "Bearer "+token)
|
req.Header.Set("Authorization", "Bearer "+token)
|
||||||
@@ -198,7 +199,7 @@ func TestRunner_ExternalCacheServer_RegisterRevoke(t *testing.T) {
|
|||||||
"token must be accepted after registerCacheForTask")
|
"token must be accepted after registerCacheForTask")
|
||||||
|
|
||||||
// The server took the results service over, so the artifact half reaches Gitea through it.
|
// The server took the results service over, so the artifact half reaches Gitea through it.
|
||||||
require.Equal(t, remote.ExternalURL(), resultsURL)
|
require.Equal(t, external, resultsURL)
|
||||||
artifact, err := http.NewRequestWithContext(t.Context(), http.MethodPost,
|
artifact, err := http.NewRequestWithContext(t.Context(), http.MethodPost,
|
||||||
resultsURL+"/twirp/github.actions.results.api.v1.ArtifactService/CreateArtifact", nil)
|
resultsURL+"/twirp/github.actions.results.api.v1.ArtifactService/CreateArtifact", nil)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
@@ -248,7 +249,7 @@ func TestRunner_ExternalCacheServer_RegisterRevoke(t *testing.T) {
|
|||||||
ArchiveLocation string `json:"archiveLocation"`
|
ArchiveLocation string `json:"archiveLocation"`
|
||||||
}
|
}
|
||||||
require.NoError(t, decodeJSON(resp, &hit))
|
require.NoError(t, decodeJSON(resp, &hit))
|
||||||
require.NotEmpty(t, hit.ArchiveLocation)
|
require.True(t, strings.HasPrefix(hit.ArchiveLocation, external), hit.ArchiveLocation)
|
||||||
|
|
||||||
dl, err := http.Get(hit.ArchiveLocation)
|
dl, err := http.Get(hit.ArchiveLocation)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|||||||
@@ -137,6 +137,7 @@ cache:
|
|||||||
#port: 0
|
#port: 0
|
||||||
# URL of a shared `gitea-runner cache-server` to use instead of starting a local one.
|
# URL of a shared `gitea-runner cache-server` to use instead of starting a local one.
|
||||||
# Set on every runner that should share a cache pool. A trailing slash is optional.
|
# Set on every runner that should share a cache pool. A trailing slash is optional.
|
||||||
|
# Jobs reach the server at this URL too, so set it to the reverse proxy when one fronts the server.
|
||||||
# Example: "http://cache-host:8088/"
|
# Example: "http://cache-host:8088/"
|
||||||
# Requires external_secret (below) to match the value on the cache-server.
|
# Requires external_secret (below) to match the value on the cache-server.
|
||||||
#external_server: ""
|
#external_server: ""
|
||||||
|
|||||||
Reference in New Issue
Block a user