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
}
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{
Name: key,
Durable: key,
Name: consumerName,
Durable: consumerName,
FilterSubject: strings.TrimSuffix(subjectBase, ".") + "." + key,
AckPolicy: jetstream.AckExplicitPolicy,
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 {
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)
}
}
+20 -13
View File
@@ -13,13 +13,12 @@ import (
)
const (
EnvAssignmentID = "CI_ASSIGNMENT_ID"
EnvCapability = "CI_RUNNER_CAPABILITY"
EnvFacadeURL = "CI_RUNNER_FACADE_URL"
EnvFacadeID = "CI_RUNNER_FACADE_SPIFFE_ID"
EnvSPIFFEID = "CI_SPIFFE_ID"
EnvWorkloadClass = "CI_WORKLOAD_CLASS"
EnvDriver = "CI_WORKLOAD_DRIVER"
EnvAssignmentID = "CI_ASSIGNMENT_ID"
EnvCapability = "CI_RUNNER_CAPABILITY"
EnvFacadeURL = "CI_RUNNER_FACADE_URL"
EnvFacadeID = "CI_RUNNER_FACADE_SPIFFE_ID"
EnvSPIFFEID = "CI_SPIFFE_ID"
EnvRunnerLabels = "CI_RUNNER_LABELS_JSON"
)
// Bootstrap emits assignment-scoped launch configuration. FacadeURL is the
@@ -43,14 +42,17 @@ func (b Bootstrap) Environment(assignment taskassignment.Assignment) (map[string
if capability == "" {
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{
EnvAssignmentID: assignment.ID,
EnvCapability: capability,
EnvFacadeURL: b.FacadeURL,
EnvFacadeID: b.FacadeSPIFFEID,
EnvSPIFFEID: assignment.Identity.SPIFFEID,
EnvWorkloadClass: string(assignment.Placement.Class),
EnvDriver: string(assignment.Placement.Driver),
EnvRunnerLabels: string(labels),
"SPIFFE_ENDPOINT_SOCKET": b.WorkloadAPIAddr,
}, nil
}
@@ -69,7 +71,7 @@ type Registration struct {
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 == "" {
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 == "" {
return nil, errors.New("local runner proxy URL must be an absolute http URL")
}
if err := placement.Validate(); err != nil {
return nil, err
if len(labels) == 0 {
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{
Warning: "Generated for one preassigned task by gitea-dynamic-runner.",
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, "", " ")
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 {
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)
}
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) {
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 {
t.Fatal(err)
}
@@ -72,7 +72,7 @@ func TestRegistrationMatchesOfficialRunnerSchema(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")
}
}
+7 -8
View File
@@ -2,6 +2,7 @@ package runnerbootstrap
import (
"context"
"encoding/json"
"errors"
"fmt"
"net"
@@ -10,14 +11,12 @@ import (
"os/exec"
"path/filepath"
"time"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskassignment"
)
type ExecutorConfig struct {
AssignmentID string
Capability string
Placement taskassignment.Placement
RunnerLabels []string
FacadeURL string
FacadeSPIFFEID string
WorkloadAPIAddr string
@@ -71,7 +70,7 @@ func RunExecutor(ctx context.Context, config ExecutorConfig) error {
defer os.RemoveAll(workDir)
}
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 {
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
// address follows SPIFFE_ENDPOINT_SOCKET through go-spiffe when not set here.
func ExecutorConfigFromEnvironment() (ExecutorConfig, error) {
placement := taskassignment.Placement{Class: taskassignment.WorkloadClass(os.Getenv(EnvWorkloadClass)), Driver: taskassignment.Driver(os.Getenv(EnvDriver))}
if err := placement.Validate(); err != nil {
return ExecutorConfig{}, err
var labels []string
if err := json.Unmarshal([]byte(os.Getenv(EnvRunnerLabels)), &labels); err != nil || len(labels) == 0 {
return ExecutorConfig{}, errors.New("valid runner labels are required")
}
config := ExecutorConfig{
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"),
ListenAddress: "127.0.0.1:0",
Stdout: os.Stdout, Stderr: os.Stderr,