Files
gitea-dynamic-runner/internal/runnerfacade/facade.go

146 lines
5.7 KiB
Go

// Package runnerfacade presents pre-assigned tasks to unmodified Gitea Runner binaries.
package runnerfacade
import (
"context"
"errors"
"net/http"
"connectrpc.com/connect"
"gitea.dev/actionslib/pkg/protocol"
runnerv1 "gitea.dev/actionslib/runner/v1"
"gitea.dev/actionslib/runner/v1/runnerv1connect"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskassignment"
)
type identityKey struct{}
func WithSPIFFEID(ctx context.Context, id string) context.Context {
return context.WithValue(ctx, identityKey{}, id)
}
// SPIFFEMiddleware extracts the authenticated workload identity from the mTLS
// peer certificate. TLS verification itself is configured by the server with
// go-spiffe; this layer only passes the verified ID into Connect handlers.
func SPIFFEMiddleware(next http.Handler) http.Handler {
return http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) {
if request.TLS == nil || len(request.TLS.PeerCertificates) == 0 {
http.Error(response, "client SPIFFE identity required", http.StatusUnauthorized)
return
}
var spiffeID string
for _, uri := range request.TLS.PeerCertificates[0].URIs {
if uri.Scheme == "spiffe" {
if spiffeID != "" {
http.Error(response, "multiple client SPIFFE identities", http.StatusUnauthorized)
return
}
spiffeID = uri.String()
}
}
if spiffeID == "" {
http.Error(response, "client SPIFFE identity required", http.StatusUnauthorized)
return
}
next.ServeHTTP(response, request.WithContext(WithSPIFFEID(request.Context(), spiffeID)))
})
}
type Upstream interface {
UpdateTask(context.Context, *connect.Request[runnerv1.UpdateTaskRequest]) (*connect.Response[runnerv1.UpdateTaskResponse], error)
UpdateLog(context.Context, *connect.Request[runnerv1.UpdateLogRequest]) (*connect.Response[runnerv1.UpdateLogResponse], error)
}
type Facade struct {
runnerv1connect.UnimplementedRunnerServiceHandler
Registry *Registry
Capabilities Capabilities
Upstream Upstream
OnTerminal func(context.Context, taskassignment.Assignment) error
}
func (f *Facade) Handler() (string, http.Handler) {
return runnerv1connect.NewRunnerServiceHandler(f)
}
func (f *Facade) Register(context.Context, *connect.Request[runnerv1.RegisterRequest]) (*connect.Response[runnerv1.RegisterResponse], error) {
return nil, connect.NewError(connect.CodePermissionDenied, errors.New("executor registration is disabled"))
}
func (f *Facade) Declare(ctx context.Context, request *connect.Request[runnerv1.DeclareRequest]) (*connect.Response[runnerv1.DeclareResponse], error) {
assignmentID, _, err := f.authenticate(ctx, request)
if err != nil {
return nil, err
}
return connect.NewResponse(&runnerv1.DeclareResponse{Runner: &runnerv1.Runner{
Uuid: assignmentID, Name: assignmentID, Status: runnerv1.RunnerStatus_RUNNER_STATUS_IDLE,
Version: request.Msg.GetVersion(), Labels: append([]string(nil), request.Msg.GetLabels()...), Ephemeral: true,
}}), nil
}
func (f *Facade) FetchTask(ctx context.Context, request *connect.Request[runnerv1.FetchTaskRequest]) (*connect.Response[runnerv1.FetchTaskResponse], error) {
assignmentID, spiffeID, err := f.authenticate(ctx, request)
if err != nil {
return nil, err
}
assignment, err := f.Registry.Claim(assignmentID, spiffeID)
if err != nil {
return nil, connect.NewError(connect.CodeFailedPrecondition, err)
}
return connect.NewResponse(&runnerv1.FetchTaskResponse{Task: assignment.Task}), nil
}
func (f *Facade) UpdateTask(ctx context.Context, request *connect.Request[runnerv1.UpdateTaskRequest]) (*connect.Response[runnerv1.UpdateTaskResponse], error) {
assignment, err := f.authorizeClaimed(ctx, request)
if err != nil {
return nil, err
}
if request.Msg.GetState().GetId() != assignment.Task.GetId() {
return nil, connect.NewError(connect.CodePermissionDenied, errors.New("task update does not match assignment"))
}
response, err := f.Upstream.UpdateTask(ctx, connect.NewRequest(request.Msg))
if err == nil && request.Msg.GetState().GetResult() != runnerv1.Result_RESULT_UNSPECIFIED && f.OnTerminal != nil {
if terminalErr := f.OnTerminal(ctx, assignment); terminalErr != nil {
return nil, connect.NewError(connect.CodeUnavailable, terminalErr)
}
}
return response, err
}
func (f *Facade) UpdateLog(ctx context.Context, request *connect.Request[runnerv1.UpdateLogRequest]) (*connect.Response[runnerv1.UpdateLogResponse], error) {
assignment, err := f.authorizeClaimed(ctx, request)
if err != nil {
return nil, err
}
if request.Msg.GetTaskId() != assignment.Task.GetId() {
return nil, connect.NewError(connect.CodePermissionDenied, errors.New("log update does not match assignment"))
}
return f.Upstream.UpdateLog(ctx, connect.NewRequest(request.Msg))
}
func (f *Facade) authorizeClaimed(ctx context.Context, request connect.AnyRequest) (taskassignment.Assignment, error) {
assignmentID, spiffeID, err := f.authenticate(ctx, request)
if err != nil {
return taskassignment.Assignment{}, err
}
assignment, err := f.Registry.Resolve(assignmentID, spiffeID)
if err != nil {
return taskassignment.Assignment{}, connect.NewError(connect.CodePermissionDenied, err)
}
return assignment, nil
}
func (f *Facade) authenticate(ctx context.Context, request connect.AnyRequest) (string, string, error) {
if f.Registry == nil || f.Upstream == nil {
return "", "", connect.NewError(connect.CodeInternal, errors.New("runner facade is not configured"))
}
assignmentID := request.Header().Get(protocol.UUIDHeader)
token := request.Header().Get(protocol.TokenHeader)
spiffeID, _ := ctx.Value(identityKey{}).(string)
if assignmentID == "" || spiffeID == "" || !f.Capabilities.Verify(assignmentID, token) {
return "", "", connect.NewError(connect.CodeUnauthenticated, errors.New("invalid executor identity or capability"))
}
return assignmentID, spiffeID, nil
}