Files
panxiao81 48b6b8038e
test / shell (pull_request) Successful in 21s
test / go (pull_request) Successful in 5m12s
test / python (pull_request) Failing after 11m44s
feat: 隔离 VM 集成测试标签
2026-09-21 08:50:18 +00:00

153 lines
5.0 KiB
Go

// Package taskassignment defines the durable handoff between the scheduler and workers.
package taskassignment
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"slices"
"strconv"
"gitea.dev/actionslib/pkg/model"
runnerv1 "gitea.dev/actionslib/runner/v1"
"google.golang.org/protobuf/proto"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskidentity"
)
const wireVersion = 1
type Backend string
const (
BackendPod Backend = "pod"
BackendVM Backend = "vm"
)
// Assignment is the only document persisted in the handoff queue.
type Assignment struct {
ID string
Backend Backend
Task *runnerv1.Task
Identity taskidentity.Identity
}
// FromMetadata reconstructs the minimal assignment needed to authorize an
// already-running executor after a controller restart. Backend metadata was
// originally derived from the trusted Gitea task and is validated again here.
func FromMetadata(labels, annotations map[string]string, trustDomain string) (Assignment, error) {
taskID, err := strconv.ParseInt(labels["ci.ddupan.top/task-id"], 10, 64)
if err != nil || taskID < 1 {
return Assignment{}, errors.New("backend metadata has invalid task ID")
}
backend := Backend(labels["ci.ddupan.top/backend"])
if backend != BackendPod && backend != BackendVM {
return Assignment{}, errors.New("backend metadata has invalid backend")
}
id := labels["ci.ddupan.top/assignment-id"]
if id != fmt.Sprintf("gitea-task-%d", taskID) {
return Assignment{}, errors.New("backend metadata assignment ID does not match task ID")
}
identity, err := taskidentity.FromMetadata(
annotations["ci.ddupan.top/repository"], annotations["ci.ddupan.top/job-key"],
annotations["ci.ddupan.top/spiffe-id"], trustDomain,
)
if err != nil {
return Assignment{}, err
}
return Assignment{ID: id, Backend: backend, Task: &runnerv1.Task{Id: taskID}, Identity: identity}, nil
}
type envelope struct {
Version int `json:"version"`
ID string `json:"id"`
Backend Backend `json:"backend"`
Task []byte `json:"task"`
Identity taskidentity.Identity `json:"identity"`
}
// New derives all trusted assignment fields from the task fetched from Gitea.
func New(task *runnerv1.Task, trustDomain string) (Assignment, error) {
if task == nil || task.GetId() <= 0 {
return Assignment{}, errors.New("positive Gitea task ID is required")
}
identity, err := taskidentity.FromTask(task, trustDomain)
if err != nil {
return Assignment{}, err
}
backend, err := backendFromTask(task)
if err != nil {
return Assignment{}, err
}
return Assignment{
ID: fmt.Sprintf("gitea-task-%d", task.GetId()),
Backend: backend,
Task: task,
Identity: identity,
}, nil
}
func backendFromTask(task *runnerv1.Task) (Backend, error) {
workflow, err := model.ReadWorkflow(bytes.NewReader(task.GetWorkflowPayload()))
if err != nil {
return "", fmt.Errorf("parse task workflow for backend: %w", err)
}
jobIDs := workflow.GetJobIDs()
if len(jobIDs) != 1 || workflow.GetJob(jobIDs[0]) == nil {
return "", fmt.Errorf("task workflow must contain exactly one non-empty job")
}
labels := workflow.GetJob(jobIDs[0]).RunsOnLabels()
if !slices.Contains(labels, "self-hosted") {
return "", fmt.Errorf("task runs-on labels must include self-hosted: %v", labels)
}
hasPod := slices.Contains(labels, string(BackendPod))
hasVM := slices.Contains(labels, string(BackendVM)) || slices.Contains(labels, "vm-dev")
if hasPod == hasVM {
return "", fmt.Errorf("task runs-on labels must select exactly one of pod or vm: %v", labels)
}
if hasPod {
return BackendPod, nil
}
return BackendVM, nil
}
// Marshal encodes a versioned assignment. Protobuf preserves the exact Gitea task.
func Marshal(assignment Assignment) ([]byte, error) {
if assignment.Task == nil {
return nil, errors.New("assignment task is required")
}
task, err := proto.Marshal(assignment.Task)
if err != nil {
return nil, fmt.Errorf("marshal Gitea task: %w", err)
}
return json.Marshal(envelope{
Version: wireVersion,
ID: assignment.ID, Backend: assignment.Backend,
Task: task, Identity: assignment.Identity,
})
}
// Unmarshal re-derives trusted fields instead of trusting duplicated queue metadata.
func Unmarshal(data []byte, trustDomain string) (Assignment, error) {
var wire envelope
if err := json.Unmarshal(data, &wire); err != nil {
return Assignment{}, fmt.Errorf("decode assignment: %w", err)
}
if wire.Version != wireVersion {
return Assignment{}, fmt.Errorf("unsupported assignment version %d", wire.Version)
}
task := new(runnerv1.Task)
if err := proto.Unmarshal(wire.Task, task); err != nil {
return Assignment{}, fmt.Errorf("unmarshal Gitea task: %w", err)
}
canonical, err := New(task, trustDomain)
if err != nil {
return Assignment{}, err
}
if wire.ID != canonical.ID || wire.Backend != canonical.Backend || wire.Identity != canonical.Identity {
return Assignment{}, errors.New("assignment metadata does not match its Gitea task")
}
return canonical, nil
}