Compare commits

...
Author SHA1 Message Date
panxiao81 53732699df fix: 解耦 runner 与 placement 实现
test / shell (pull_request) Successful in 3m31s
test / python (pull_request) Failing after 10m54s
test / go (pull_request) Failing after 12m12s
2026-09-25 19:03:13 +00:00
panxiao81 011c1e4c83 Merge 修复 placement consumer 名称导致的启动失败
test / shell (push) Successful in 43s
test / python (push) Successful in 1m31s
publish controller image / publish-controller (push) Successful in 5m8s
test / go (push) Successful in 5m36s
2026-09-25 18:46:30 +00:00
panxiao81 1f9577c816 fix: 使用合法的 JetStream consumer 名称
test / shell (pull_request) Successful in 52s
test / python (pull_request) Successful in 1m24s
test / go (pull_request) Successful in 3m57s
2026-09-25 18:41:47 +00:00
panxiao81 6ea78d4704 Merge pull request '重构 workload class 与 placement driver' (#54) from refactor/workload-placement into main
test / shell (push) Successful in 1m7s
test / python (push) Successful in 3m3s
publish controller image / publish-controller (push) Successful in 5m44s
test / go (push) Successful in 9m7s
Reviewed-on: #54
2026-09-25 17:11:01 +00:00
5 changed files with 37 additions and 27 deletions
+6 -2
View File
@@ -186,9 +186,13 @@ func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectB
return nil, err return nil, err
} }
key := placement.Key() key := placement.Key()
// JetStream consumer names may not contain dots even though subjects do.
// Keep the placement subject hierarchical while using a stable, legal name
// for the durable cursor.
consumerName := strings.ReplaceAll(key, ".", "-")
consumer, err := manager.CreateOrUpdateConsumer(ctx, stream, jetstream.ConsumerConfig{ consumer, err := manager.CreateOrUpdateConsumer(ctx, stream, jetstream.ConsumerConfig{
Name: key, Name: consumerName,
Durable: key, Durable: consumerName,
FilterSubject: strings.TrimSuffix(subjectBase, ".") + "." + key, FilterSubject: strings.TrimSuffix(subjectBase, ".") + "." + key,
AckPolicy: jetstream.AckExplicitPolicy, AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 5 * time.Minute, AckWait: 5 * time.Minute,
+1 -1
View File
@@ -215,7 +215,7 @@ func TestOpenConsumerUsesIndependentDurablePerPlacement(t *testing.T) {
if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.KubernetesContainer, 4); err != nil { if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.KubernetesContainer, 4); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if manager.config.Durable != "container.kubernetes" || manager.config.FilterSubject != "ci.assignment.container.kubernetes" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 { if manager.config.Name != "container-kubernetes" || manager.config.Durable != "container-kubernetes" || manager.config.FilterSubject != "ci.assignment.container.kubernetes" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 {
t.Fatalf("config = %#v", manager.config) t.Fatalf("config = %#v", manager.config)
} }
} }
+20 -13
View File
@@ -13,13 +13,12 @@ import (
) )
const ( const (
EnvAssignmentID = "CI_ASSIGNMENT_ID" EnvAssignmentID = "CI_ASSIGNMENT_ID"
EnvCapability = "CI_RUNNER_CAPABILITY" EnvCapability = "CI_RUNNER_CAPABILITY"
EnvFacadeURL = "CI_RUNNER_FACADE_URL" EnvFacadeURL = "CI_RUNNER_FACADE_URL"
EnvFacadeID = "CI_RUNNER_FACADE_SPIFFE_ID" EnvFacadeID = "CI_RUNNER_FACADE_SPIFFE_ID"
EnvSPIFFEID = "CI_SPIFFE_ID" EnvSPIFFEID = "CI_SPIFFE_ID"
EnvWorkloadClass = "CI_WORKLOAD_CLASS" EnvRunnerLabels = "CI_RUNNER_LABELS_JSON"
EnvDriver = "CI_WORKLOAD_DRIVER"
) )
// Bootstrap emits assignment-scoped launch configuration. FacadeURL is the // Bootstrap emits assignment-scoped launch configuration. FacadeURL is the
@@ -43,14 +42,17 @@ func (b Bootstrap) Environment(assignment taskassignment.Assignment) (map[string
if capability == "" { if capability == "" {
return nil, errors.New("runner capability issuer is not configured") return nil, errors.New("runner capability issuer is not configured")
} }
labels, err := json.Marshal([]string{"self-hosted", string(assignment.Placement.Class), string(assignment.Placement.Driver)})
if err != nil {
return nil, fmt.Errorf("encode runner labels: %w", err)
}
return map[string]string{ return map[string]string{
EnvAssignmentID: assignment.ID, EnvAssignmentID: assignment.ID,
EnvCapability: capability, EnvCapability: capability,
EnvFacadeURL: b.FacadeURL, EnvFacadeURL: b.FacadeURL,
EnvFacadeID: b.FacadeSPIFFEID, EnvFacadeID: b.FacadeSPIFFEID,
EnvSPIFFEID: assignment.Identity.SPIFFEID, EnvSPIFFEID: assignment.Identity.SPIFFEID,
EnvWorkloadClass: string(assignment.Placement.Class), EnvRunnerLabels: string(labels),
EnvDriver: string(assignment.Placement.Driver),
"SPIFFE_ENDPOINT_SOCKET": b.WorkloadAPIAddr, "SPIFFE_ENDPOINT_SOCKET": b.WorkloadAPIAddr,
}, nil }, nil
} }
@@ -69,7 +71,7 @@ type Registration struct {
Ephemeral bool `json:"ephemeral"` Ephemeral bool `json:"ephemeral"`
} }
func RegistrationJSON(assignmentID, capability, localProxyURL string, placement taskassignment.Placement) ([]byte, error) { func RegistrationJSON(assignmentID, capability, localProxyURL string, labels []string) ([]byte, error) {
if assignmentID == "" || capability == "" { if assignmentID == "" || capability == "" {
return nil, errors.New("assignment ID and runner capability are required") return nil, errors.New("assignment ID and runner capability are required")
} }
@@ -77,13 +79,18 @@ func RegistrationJSON(assignmentID, capability, localProxyURL string, placement
if err != nil || parsed.Scheme != "http" || parsed.Host == "" { if err != nil || parsed.Scheme != "http" || parsed.Host == "" {
return nil, errors.New("local runner proxy URL must be an absolute http URL") return nil, errors.New("local runner proxy URL must be an absolute http URL")
} }
if err := placement.Validate(); err != nil { if len(labels) == 0 {
return nil, err return nil, errors.New("runner labels are required")
}
for _, label := range labels {
if label == "" {
return nil, errors.New("runner labels must not be empty")
}
} }
registration := Registration{ registration := Registration{
Warning: "Generated for one preassigned task by gitea-dynamic-runner.", Warning: "Generated for one preassigned task by gitea-dynamic-runner.",
UUID: assignmentID, Name: assignmentID, Token: capability, UUID: assignmentID, Name: assignmentID, Token: capability,
Address: localProxyURL, Labels: []string{"self-hosted", string(placement.Class), string(placement.Driver)}, Ephemeral: true, Address: localProxyURL, Labels: append([]string(nil), labels...), Ephemeral: true,
} }
data, err := json.MarshalIndent(registration, "", " ") data, err := json.MarshalIndent(registration, "", " ")
if err != nil { if err != nil {
+3 -3
View File
@@ -46,7 +46,7 @@ func TestEnvironmentIsDeterministicAndAssignmentScoped(t *testing.T) {
if first[EnvAssignmentID] != "gitea-task-42" || first[EnvSPIFFEID] != testAssignment().Identity.SPIFFEID { if first[EnvAssignmentID] != "gitea-task-42" || first[EnvSPIFFEID] != testAssignment().Identity.SPIFFEID {
t.Fatalf("environment = %#v", first) t.Fatalf("environment = %#v", first)
} }
if first[EnvWorkloadClass] != "container" || first[EnvDriver] != "kubernetes" || first[EnvFacadeID] == "" { if first[EnvRunnerLabels] != `["self-hosted","container","kubernetes"]` || first[EnvFacadeID] == "" {
t.Fatalf("environment = %#v", first) t.Fatalf("environment = %#v", first)
} }
if first["SPIFFE_ENDPOINT_SOCKET"] != "unix:///run/spire/agent-sockets/spire-agent.sock" { if first["SPIFFE_ENDPOINT_SOCKET"] != "unix:///run/spire/agent-sockets/spire-agent.sock" {
@@ -55,7 +55,7 @@ func TestEnvironmentIsDeterministicAndAssignmentScoped(t *testing.T) {
} }
func TestRegistrationMatchesOfficialRunnerSchema(t *testing.T) { func TestRegistrationMatchesOfficialRunnerSchema(t *testing.T) {
data, err := RegistrationJSON("gitea-task-42", "capability", "http://127.0.0.1:8080", taskassignment.OpenSandboxVM) data, err := RegistrationJSON("gitea-task-42", "capability", "http://127.0.0.1:8080", []string{"self-hosted", "vm", "opensandbox"})
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -72,7 +72,7 @@ func TestRegistrationMatchesOfficialRunnerSchema(t *testing.T) {
} }
func TestRegistrationRejectsNonLocalTLSAddress(t *testing.T) { func TestRegistrationRejectsNonLocalTLSAddress(t *testing.T) {
if _, err := RegistrationJSON("id", "capability", "https://facade.example", taskassignment.KubernetesContainer); err == nil { if _, err := RegistrationJSON("id", "capability", "https://facade.example", []string{"self-hosted"}); err == nil {
t.Fatal("expected local proxy URL validation error") t.Fatal("expected local proxy URL validation error")
} }
} }
+7 -8
View File
@@ -2,6 +2,7 @@ package runnerbootstrap
import ( import (
"context" "context"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"net" "net"
@@ -10,14 +11,12 @@ import (
"os/exec" "os/exec"
"path/filepath" "path/filepath"
"time" "time"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskassignment"
) )
type ExecutorConfig struct { type ExecutorConfig struct {
AssignmentID string AssignmentID string
Capability string Capability string
Placement taskassignment.Placement RunnerLabels []string
FacadeURL string FacadeURL string
FacadeSPIFFEID string FacadeSPIFFEID string
WorkloadAPIAddr string WorkloadAPIAddr string
@@ -71,7 +70,7 @@ func RunExecutor(ctx context.Context, config ExecutorConfig) error {
defer os.RemoveAll(workDir) defer os.RemoveAll(workDir)
} }
registration, err := RegistrationJSON( registration, err := RegistrationJSON(
config.AssignmentID, config.Capability, "http://"+listener.Addr().String(), config.Placement, config.AssignmentID, config.Capability, "http://"+listener.Addr().String(), config.RunnerLabels,
) )
if err != nil { if err != nil {
return err return err
@@ -144,13 +143,13 @@ func waitForFacade(ctx context.Context, endpoint string) error {
// the assignment-scoped values injected by the backend. The Workload API // the assignment-scoped values injected by the backend. The Workload API
// address follows SPIFFE_ENDPOINT_SOCKET through go-spiffe when not set here. // address follows SPIFFE_ENDPOINT_SOCKET through go-spiffe when not set here.
func ExecutorConfigFromEnvironment() (ExecutorConfig, error) { func ExecutorConfigFromEnvironment() (ExecutorConfig, error) {
placement := taskassignment.Placement{Class: taskassignment.WorkloadClass(os.Getenv(EnvWorkloadClass)), Driver: taskassignment.Driver(os.Getenv(EnvDriver))} var labels []string
if err := placement.Validate(); err != nil { if err := json.Unmarshal([]byte(os.Getenv(EnvRunnerLabels)), &labels); err != nil || len(labels) == 0 {
return ExecutorConfig{}, err return ExecutorConfig{}, errors.New("valid runner labels are required")
} }
config := ExecutorConfig{ config := ExecutorConfig{
AssignmentID: os.Getenv(EnvAssignmentID), Capability: os.Getenv(EnvCapability), AssignmentID: os.Getenv(EnvAssignmentID), Capability: os.Getenv(EnvCapability),
Placement: placement, FacadeURL: os.Getenv(EnvFacadeURL), FacadeSPIFFEID: os.Getenv(EnvFacadeID), RunnerLabels: labels, FacadeURL: os.Getenv(EnvFacadeURL), FacadeSPIFFEID: os.Getenv(EnvFacadeID),
RunnerBinary: os.Getenv("GITEA_RUNNER_BINARY"), RunnerConfig: os.Getenv("GITEA_RUNNER_CONFIG_FILE"), RunnerBinary: os.Getenv("GITEA_RUNNER_BINARY"), RunnerConfig: os.Getenv("GITEA_RUNNER_CONFIG_FILE"),
ListenAddress: "127.0.0.1:0", ListenAddress: "127.0.0.1:0",
Stdout: os.Stdout, Stderr: os.Stderr, Stdout: os.Stdout, Stderr: os.Stderr,