package httpserver import ( "context" "encoding/json" "errors" "fmt" "net" "net/http" "net/url" "os" "strings" "time" "github.com/toki/oto/services/core/internal/cicdstate" "github.com/toki/oto/services/core/internal/runnerregistry" otopb "github.com/toki/oto/services/core/oto" ) // Server wraps the HTTP server for the OTO Core service. type Server struct { httpServer *http.Server } // NewServer creates a new instance of Server. func NewServer(addr string) *Server { return NewServerWithRegistryAndStore(addr, runnerregistry.New(), cicdstate.NewStore()) } // NewServerWithRegistry creates a server using an injected runner registry. func NewServerWithRegistry(addr string, registry *runnerregistry.Registry) *Server { return NewServerWithRegistryAndStore(addr, registry, cicdstate.NewStore()) } // NewServerWithRegistryAndStore creates a server with injected runner registry and CICD store. func NewServerWithRegistryAndStore(addr string, registry *runnerregistry.Registry, store *cicdstate.Store) *Server { mux := http.NewServeMux() mux.HandleFunc("/healthz", handleHealthz) mux.HandleFunc("/readyz", handleReadyz) mux.HandleFunc("/api/v1/runners/register", handleRunnerRegister(registry)) mux.HandleFunc("/api/v1/runners/bootstrap-command", handleRunnerBootstrapCommand(registry)) mux.HandleFunc("/api/v1/runners/{id}/heartbeat", handleRunnerHeartbeat(registry)) mux.HandleFunc("/api/v1/runners/{id}/disconnect", handleRunnerDisconnect(registry)) mux.HandleFunc("/api/v1/runners/{id}", handleGetRunner(registry)) mux.HandleFunc("/bootstrap/oto-agent.sh", handleServeBootstrapScript()) mux.HandleFunc("/api/v1/", handleRouter(store, registry)) return &Server{ httpServer: &http.Server{ Addr: addr, Handler: mux, }, } } // Start starts the HTTP server. func (s *Server) Start() error { return s.httpServer.ListenAndServe() } // StartListener starts the HTTP server using a custom net.Listener. // Useful for testing with dynamic ports. func (s *Server) StartListener(ln net.Listener) error { return s.httpServer.Serve(ln) } // Shutdown gracefully shuts down the server. func (s *Server) Shutdown(ctx context.Context) error { return s.httpServer.Shutdown(ctx) } func handleHealthz(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } w.WriteHeader(http.StatusOK) w.Write([]byte("OK")) } func handleReadyz(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } w.WriteHeader(http.StatusOK) w.Write([]byte("OK")) } func handleRunnerRegister(registry *runnerregistry.Registry) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } var request otopb.RegisterRunnerRequest if err := json.NewDecoder(r.Body).Decode(&request); err != nil { writeRunnerRegisterResponse(w, http.StatusBadRequest, &otopb.RegisterRunnerResponse{ Accepted: false, RejectReason: "invalid registration request", Error: &otopb.ProtocolError{ Code: "invalid_registration_request", Message: "invalid registration request", }, }) return } response := registry.Register(&request) writeRunnerRegisterResponse(w, http.StatusOK, response) } } func handleRunnerHeartbeat(registry *runnerregistry.Registry) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } runnerID := r.PathValue("id") if runnerID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { runnerID = parts[4] } } if runnerID == "" { writeHeartbeatResponse(w, http.StatusBadRequest, &otopb.HeartbeatResponse{ Success: false, ErrorMessage: "missing runner id", }) return } var req otopb.HeartbeatRequest // Decode request if body is not empty if r.ContentLength > 0 || r.Body != nil { if err := json.NewDecoder(r.Body).Decode(&req); err != nil && err.Error() != "EOF" { writeHeartbeatResponse(w, http.StatusBadRequest, &otopb.HeartbeatResponse{ Success: false, ErrorMessage: "invalid heartbeat request", }) return } } if req.RunnerId == "" { req.RunnerId = runnerID } else if req.RunnerId != runnerID { writeHeartbeatResponse(w, http.StatusBadRequest, &otopb.HeartbeatResponse{ Success: false, ErrorMessage: "runner id mismatch between path and body", }) return } response := registry.Heartbeat(req.RunnerId, req.Status) if !response.Success { status := http.StatusNotFound if response.ErrorMessage != "unknown runner" { status = http.StatusBadRequest } writeHeartbeatResponse(w, status, response) return } writeHeartbeatResponse(w, http.StatusOK, response) } } func handleRunnerDisconnect(registry *runnerregistry.Registry) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } runnerID := r.PathValue("id") if runnerID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { runnerID = parts[4] } } if runnerID == "" { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "success": false, "error": "missing runner id", }) return } success := registry.Disconnect(runnerID) w.Header().Set("Content-Type", "application/json") if !success { w.WriteHeader(http.StatusNotFound) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "success": false, "error": "runner not found", }) return } w.WriteHeader(http.StatusOK) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "success": true, }) } } func writeRunnerRegisterResponse(w http.ResponseWriter, status int, response *otopb.RegisterRunnerResponse) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(status) _ = json.NewEncoder(w).Encode(response) } func writeHeartbeatResponse(w http.ResponseWriter, status int, response *otopb.HeartbeatResponse) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(status) _ = json.NewEncoder(w).Encode(response) } func handleGetRunner(registry *runnerregistry.Registry) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } runnerID := r.PathValue("id") if runnerID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { runnerID = parts[4] } } if runnerID == "" { http.Error(w, "missing runner id", http.StatusBadRequest) return } record, ok := registry.Snapshot(runnerID) if !ok { http.Error(w, "runner not found", http.StatusNotFound) return } w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusOK) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "runner_id": record.RunnerID, "alias": record.Alias, "protocol_version": record.ProtocolVersion, "status": record.Status, "accepted_at": record.AcceptedAt.Format(time.RFC3339), "first_heartbeat_at": record.FirstHeartbeatAt.Format(time.RFC3339), "last_heartbeat_at": record.LastHeartbeatAt.Format(time.RFC3339), "failure_reason": record.FailureReason, }) } } func shellEscape(s string) string { return "'" + strings.ReplaceAll(s, "'", "'\\''") + "'" } func handleRunnerBootstrapCommand(registry *runnerregistry.Registry) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } var request otopb.BootstrapCommandRequest if err := json.NewDecoder(r.Body).Decode(&request); err != nil { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "error": "invalid bootstrap command request", }) return } runnerID := strings.TrimSpace(request.GetRunnerId()) token := strings.TrimSpace(request.GetEnrollmentToken()) if runnerID == "" || token == "" { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "error": "missing runner id or enrollment token", }) return } scheme := "http" if r.TLS != nil { scheme = "https" } serverURL := scheme + "://" + r.Host u, err := url.Parse(serverURL) if err != nil || u.Host == "" || u.Scheme == "" { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "error": "invalid server URL", }) return } // Host validation to prevent command injection and host header spoofing for _, char := range u.Host { if !((char >= 'a' && char <= 'z') || (char >= 'A' && char <= 'Z') || (char >= '0' && char <= '9') || char == '.' || char == '-' || char == ':' || char == '[' || char == ']') { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "error": "invalid characters in server Host", }) return } } releaseBaseURL := os.Getenv("OTO_RUNNER_RELEASE_BASE_URL") if releaseBaseURL == "" { if scheme == "https" { releaseBaseURL = serverURL + "/releases" } else { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "error": "release base URL must use HTTPS (set OTO_RUNNER_RELEASE_BASE_URL environment variable)", }) return } } else { if !strings.HasPrefix(releaseBaseURL, "https://") { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusBadRequest) _ = json.NewEncoder(w).Encode(map[string]interface{}{ "error": "OTO_RUNNER_RELEASE_BASE_URL must use https:// scheme", }) return } } escapedScriptURL := shellEscape(serverURL + "/bootstrap/oto-agent.sh") escapedServerURL := shellEscape(serverURL) escapedRunnerID := shellEscape(runnerID) escapedToken := shellEscape(token) escapedReleaseURL := shellEscape(releaseBaseURL) bootstrapCmd := fmt.Sprintf( "curl -fsSL %s | bash -s -- --server-url %s --agent-id %s --enrollment-token %s --release-base-url %s", escapedScriptURL, escapedServerURL, escapedRunnerID, escapedToken, escapedReleaseURL, ) w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusOK) _ = json.NewEncoder(w).Encode(&otopb.BootstrapCommandResponse{ BootstrapCommand: bootstrapCmd, }) } } func handleServeBootstrapScript() http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed) return } paths := []string{ "../../apps/runner/assets/script/shell/oto_agent_bootstrap.sh", "../../../../apps/runner/assets/script/shell/oto_agent_bootstrap.sh", "apps/runner/assets/script/shell/oto_agent_bootstrap.sh", } var content []byte var err error for _, p := range paths { content, err = os.ReadFile(p) if err == nil { break } } if err != nil { http.Error(w, "Bootstrap script not found: "+err.Error(), http.StatusNotFound) return } w.Header().Set("Content-Type", "application/x-sh") w.WriteHeader(http.StatusOK) _, _ = w.Write(content) } } // Extract path segments after "/api/v1/" func apiPathSegments(path string) []string { if !strings.HasPrefix(path, "/api/v1/") { return nil } trimmed := strings.TrimPrefix(path, "/api/v1/") trimmed = strings.TrimRight(trimmed, "/") if trimmed == "" { return nil } return strings.Split(trimmed, "/") } func handleRouter(store *cicdstate.Store, registry *runnerregistry.Registry) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { parts := apiPathSegments(r.URL.Path) if parts == nil { http.NotFound(w, r) return } switch parts[0] { case "jobs": switch { case len(parts) == 1 && r.Method == http.MethodPost: handleCreateJob(store)(w, r) case len(parts) == 1 && r.Method == http.MethodGet: writeResponse(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"}) return case len(parts) == 2 && r.Method == http.MethodGet: handleGetJob(store)(w, r) case len(parts) == 3 && parts[2] == "executions" && r.Method == http.MethodPost: handleCreateExecution(store)(w, r) default: http.NotFound(w, r) } case "executions": if len(parts) < 2 { http.NotFound(w, r) return } switch { case len(parts) == 2 && r.Method == http.MethodGet: handleGetExecution(store)(w, r) case len(parts) == 3 && parts[2] == "logs" && r.Method == http.MethodPost: handleAppendLog(store)(w, r) case len(parts) == 3 && parts[2] == "logs" && r.Method == http.MethodGet: handleGetLogs(store)(w, r) case len(parts) == 3 && parts[2] == "artifacts" && r.Method == http.MethodPost: handleAppendArtifact(store)(w, r) case len(parts) == 3 && parts[2] == "artifacts" && r.Method == http.MethodGet: handleGetArtifacts(store)(w, r) default: http.NotFound(w, r) } case "runners": if len(parts) < 2 { http.NotFound(w, r) return } runnerID := parts[1] switch { case len(parts) == 3 && parts[2] == "status" && r.Method == http.MethodGet: handleRunnerStatus(store, registry, runnerID)(w, r) case len(parts) == 3 && parts[2] == "self-update" && r.Method == http.MethodPost: handleRunnerSelfUpdate(store, registry, runnerID)(w, r) case len(parts) == 4 && parts[2] == "jobs" && parts[3] == "claim" && r.Method == http.MethodPost: handleRunnerClaimJob(store, registry, runnerID)(w, r) case len(parts) == 5 && parts[2] == "executions" && parts[4] == "cancel" && r.Method == http.MethodPost: handleRunnerCancelExecution(store, registry, runnerID, parts[3])(w, r) case len(parts) == 5 && parts[2] == "executions" && parts[4] == "report" && r.Method == http.MethodPost: handleRunnerReportExecution(store, registry, runnerID, parts[3])(w, r) case len(parts) == 5 && parts[2] == "executions" && parts[4] == "logs" && r.Method == http.MethodPost: handleRunnerAppendLog(store, registry, runnerID, parts[3])(w, r) case len(parts) == 5 && parts[2] == "executions" && parts[4] == "artifacts" && r.Method == http.MethodPost: handleRunnerAppendArtifact(store, registry, runnerID, parts[3])(w, r) default: http.NotFound(w, r) } default: http.NotFound(w, r) } } } func handleCreateJob(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req struct { ID string `json:"id"` Name string `json:"name"` RunRequest *otopb.RunRequest `json:"run_request"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "invalid request body"}) return } if req.ID == "" || req.Name == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "id and name are required"}) return } job, err := store.CreateJob(req.ID, req.Name, runRequestFromJSON(req.RunRequest)) if err != nil { writeResponse(w, http.StatusConflict, map[string]string{"error": err.Error()}) return } writeResponse(w, http.StatusCreated, jobToJSON(job)) } } func handleGetJob(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { jobID := r.PathValue("id") if jobID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { jobID = parts[4] } } if jobID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "missing job id"}) return } job, err := store.GetJob(jobID) if err != nil { if errors.Is(err, fmt.Errorf("job not found: %s", jobID)) || err.Error() == "job not found: "+jobID { writeResponse(w, http.StatusNotFound, map[string]string{"error": "job not found"}) return } writeResponse(w, http.StatusInternalServerError, map[string]string{"error": "internal error"}) return } writeResponse(w, http.StatusOK, jobToJSON(job)) } } func handleCreateExecution(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { jobID := r.PathValue("id") if jobID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { jobID = parts[4] } } if jobID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "missing job id"}) return } var req struct { ID string `json:"id"` } _ = json.NewDecoder(r.Body).Decode(&req) if req.ID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "execution id is required"}) return } exec, err := store.CreateExecution(jobID, req.ID) if err != nil { writeResponse(w, http.StatusBadRequest, map[string]string{"error": err.Error()}) return } writeResponse(w, http.StatusCreated, execToJSON(exec)) } } func handleGetExecution(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { execID := r.PathValue("id") if execID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { execID = parts[4] } } if execID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "missing execution id"}) return } exec, err := store.GetExecution(execID) if err != nil { writeResponse(w, http.StatusNotFound, map[string]string{"error": "execution not found"}) return } writeResponse(w, http.StatusOK, execToJSON(exec)) } } func handleAppendLog(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { writeResponse(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"}) return } execID := r.PathValue("id") if execID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { execID = parts[4] } } if execID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "missing execution id"}) return } var req struct { Line string `json:"line"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil || req.Line == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "line is required"}) return } if err := store.AppendLog(execID, req.Line); err != nil { writeResponse(w, http.StatusNotFound, map[string]string{"error": "execution not found"}) return } writeResponse(w, http.StatusCreated, map[string]string{"status": "accepted"}) } } func handleGetLogs(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { execID := r.PathValue("id") if execID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { execID = parts[4] } } if execID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "missing execution id"}) return } logs, err := store.GetLogs(execID) if err != nil { writeResponse(w, http.StatusNotFound, map[string]string{"error": "execution not found"}) return } writeResponse(w, http.StatusOK, logsToJSON(logs)) } } func handleAppendArtifact(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { writeResponse(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"}) return } execID := r.PathValue("id") if execID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { execID = parts[4] } } if execID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "missing execution id"}) return } var req struct { Name string `json:"name"` Path string `json:"path"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil || req.Name == "" || req.Path == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "name and path are required"}) return } if err := store.AppendArtifact(execID, req.Name, req.Path); err != nil { writeResponse(w, http.StatusNotFound, map[string]string{"error": "execution not found"}) return } writeResponse(w, http.StatusCreated, map[string]string{"status": "accepted"}) } } func handleGetArtifacts(store *cicdstate.Store) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { execID := r.PathValue("id") if execID == "" { parts := strings.Split(r.URL.Path, "/") if len(parts) >= 5 { execID = parts[4] } } if execID == "" { writeResponse(w, http.StatusBadRequest, map[string]string{"error": "missing execution id"}) return } artifacts, err := store.GetArtifacts(execID) if err != nil { writeResponse(w, http.StatusNotFound, map[string]string{"error": "execution not found"}) return } writeResponse(w, http.StatusOK, artifactsToJSON(artifacts)) } } func handleRunnerClaimJob(store *cicdstate.Store, registry *runnerregistry.Registry, runnerID string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req otopb.JobClaimRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.JobClaimResponse{ ErrorMessage: "invalid claim request", }) return } if req.RunnerId == "" { req.RunnerId = runnerID } else if req.RunnerId != runnerID { writeResponse(w, http.StatusBadRequest, &otopb.JobClaimResponse{ ErrorMessage: "runner id mismatch between path and body", }) return } if err := ensureRunnerKnown(registry, runnerID); err != nil { writeResponse(w, http.StatusNotFound, &otopb.JobClaimResponse{ RunnerId: runnerID, ErrorMessage: err.Error(), }) return } jobID := strings.TrimSpace(req.GetJobId()) execID := strings.TrimSpace(req.GetExecutionId()) if jobID == "" && execID == "" { var foundJob *cicdstate.Job for _, j := range store.SnapshotJobs() { if j.State == cicdstate.StateQueued { if foundJob == nil || j.CreatedAt.Before(foundJob.CreatedAt) { copyJ := j foundJob = ©J } } } if foundJob == nil { writeResponse(w, http.StatusOK, &otopb.JobClaimResponse{ Accepted: false, RunnerId: runnerID, }) return } jobID = foundJob.ID execID = fmt.Sprintf("exec-%d", time.Now().UnixNano()) } else if jobID == "" || execID == "" { writeResponse(w, http.StatusBadRequest, &otopb.JobClaimResponse{ RunnerId: runnerID, ErrorMessage: "job id and execution id are required", }) return } job, err := store.GetJob(jobID) if err != nil { writeResponse(w, http.StatusNotFound, &otopb.JobClaimResponse{ RunnerId: runnerID, JobId: jobID, ErrorMessage: "job not found", }) return } if job.State != cicdstate.StateQueued { writeResponse(w, http.StatusConflict, &otopb.JobClaimResponse{ RunnerId: runnerID, JobId: jobID, ErrorMessage: "job is not queued", State: job.State, }) return } if job.RunInput == nil { writeResponse(w, http.StatusConflict, &otopb.JobClaimResponse{ RunnerId: runnerID, JobId: jobID, ErrorMessage: "job has no run request", State: job.State, }) return } if err := validateRunInput(job.RunInput); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.JobClaimResponse{ RunnerId: runnerID, JobId: jobID, ErrorMessage: err.Error(), }) return } if _, err := store.CreateExecution(jobID, execID); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.JobClaimResponse{ RunnerId: runnerID, JobId: jobID, ExecutionId: execID, ErrorMessage: err.Error(), }) return } _ = store.SetExecutionRunnerID(execID, runnerID) if err := store.TransitionJob(jobID, cicdstate.StateRunning); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.JobClaimResponse{ RunnerId: runnerID, JobId: jobID, ExecutionId: execID, ErrorMessage: err.Error(), }) return } if err := store.TransitionExecution(execID, cicdstate.StateRunning); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.JobClaimResponse{ RunnerId: runnerID, JobId: jobID, ExecutionId: execID, ErrorMessage: err.Error(), }) return } writeResponse(w, http.StatusOK, &otopb.JobClaimResponse{ Accepted: true, RunnerId: runnerID, JobId: jobID, ExecutionId: execID, State: cicdstate.StateRunning, RunRequest: runRequestToProto(job.RunInput, runnerID, jobID, execID), }) } } func handleRunnerReportExecution(store *cicdstate.Store, registry *runnerregistry.Registry, runnerID string, execID string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req otopb.ExecutionReportRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.ExecutionReportResponse{ ErrorMessage: "invalid execution report request", }) return } if req.RunnerId == "" { req.RunnerId = runnerID } else if req.RunnerId != runnerID { writeResponse(w, http.StatusBadRequest, &otopb.ExecutionReportResponse{ ErrorMessage: "runner id mismatch between path and body", }) return } if req.ExecutionId == "" { req.ExecutionId = execID } else if req.ExecutionId != execID { writeResponse(w, http.StatusBadRequest, &otopb.ExecutionReportResponse{ ErrorMessage: "execution id mismatch between path and body", }) return } if err := ensureRunnerKnown(registry, runnerID); err != nil { writeResponse(w, http.StatusNotFound, &otopb.ExecutionReportResponse{ RunnerId: runnerID, ExecutionId: execID, ErrorMessage: err.Error(), }) return } exec, err := store.GetExecution(execID) if err != nil { writeResponse(w, http.StatusNotFound, &otopb.ExecutionReportResponse{ RunnerId: runnerID, ExecutionId: execID, ErrorMessage: "execution not found", }) return } jobID := strings.TrimSpace(req.GetJobId()) if jobID == "" { jobID = exec.JobID } else if jobID != exec.JobID { writeResponse(w, http.StatusBadRequest, &otopb.ExecutionReportResponse{ RunnerId: runnerID, JobId: jobID, ExecutionId: execID, ErrorMessage: "job id mismatch for execution", }) return } targetState := cicdstate.StateFailed if req.GetSuccess() { targetState = cicdstate.StateSucceeded } if err := store.TransitionExecution(execID, targetState); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.ExecutionReportResponse{ RunnerId: runnerID, JobId: jobID, ExecutionId: execID, ErrorMessage: err.Error(), }) return } if err := store.TransitionJob(jobID, targetState); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.ExecutionReportResponse{ RunnerId: runnerID, JobId: jobID, ExecutionId: execID, ErrorMessage: err.Error(), }) return } if strings.TrimSpace(req.GetMessage()) != "" { _ = store.AppendLog(execID, req.GetMessage()) } writeResponse(w, http.StatusOK, &otopb.ExecutionReportResponse{ Accepted: true, RunnerId: runnerID, JobId: jobID, ExecutionId: execID, State: targetState, }) } } func handleRunnerAppendLog(store *cicdstate.Store, registry *runnerregistry.Registry, runnerID string, execID string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req otopb.LogAppendRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.LogAppendResponse{ ErrorMessage: "invalid log append request", }) return } if err := normalizeRunnerExecutionRequest(registry, runnerID, execID, req.GetRunnerId(), req.GetExecutionId()); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.LogAppendResponse{ ErrorMessage: err.Error(), }) return } line := strings.TrimSpace(req.GetLine()) if line == "" { writeResponse(w, http.StatusBadRequest, &otopb.LogAppendResponse{ ErrorMessage: "line is required", }) return } if err := store.AppendLog(execID, line); err != nil { writeResponse(w, http.StatusNotFound, &otopb.LogAppendResponse{ ErrorMessage: "execution not found", }) return } writeResponse(w, http.StatusCreated, &otopb.LogAppendResponse{ Accepted: true, }) } } func handleRunnerAppendArtifact(store *cicdstate.Store, registry *runnerregistry.Registry, runnerID string, execID string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req otopb.ArtifactReportRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.ArtifactReportResponse{ ErrorMessage: "invalid artifact report request", }) return } if err := normalizeRunnerExecutionRequest(registry, runnerID, execID, req.GetRunnerId(), req.GetExecutionId()); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.ArtifactReportResponse{ ErrorMessage: err.Error(), }) return } name := strings.TrimSpace(req.GetName()) path := strings.TrimSpace(req.GetPath()) if name == "" || path == "" { writeResponse(w, http.StatusBadRequest, &otopb.ArtifactReportResponse{ ErrorMessage: "name and path are required", }) return } if err := store.AppendArtifact(execID, name, path); err != nil { writeResponse(w, http.StatusNotFound, &otopb.ArtifactReportResponse{ ErrorMessage: "execution not found", }) return } writeResponse(w, http.StatusCreated, &otopb.ArtifactReportResponse{ Accepted: true, }) } } func normalizeRunnerExecutionRequest(registry *runnerregistry.Registry, runnerID string, execID string, bodyRunnerID string, bodyExecID string) error { if bodyRunnerID != "" && bodyRunnerID != runnerID { return fmt.Errorf("runner id mismatch between path and body") } if bodyExecID != "" && bodyExecID != execID { return fmt.Errorf("execution id mismatch between path and body") } return ensureRunnerKnown(registry, runnerID) } func ensureRunnerKnown(registry *runnerregistry.Registry, runnerID string) error { if strings.TrimSpace(runnerID) == "" { return fmt.Errorf("missing runner id") } if registry == nil { return nil } if _, ok := registry.Snapshot(runnerID); !ok { return fmt.Errorf("runner not found") } return nil } // --- JSON helpers --- type jobResponse struct { ID string `json:"id"` Name string `json:"name"` State string `json:"state"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` ExecutionID string `json:"execution_id,omitempty"` } func jobToJSON(j *cicdstate.Job) map[string]interface{} { out := map[string]interface{}{ "id": j.ID, "name": j.Name, "state": j.State, "created_at": j.CreatedAt.Format(time.RFC3339), "updated_at": j.UpdatedAt.Format(time.RFC3339), "execution_id": j.ExecutionID, } if j.RunInput != nil { out["run_request"] = runInputToJSON(j.RunInput) } return out } // validateRunInput enforces the remote claim contract: exactly one of // pipeline_yaml or pipeline_yaml_path must be provided. func validateRunInput(in *cicdstate.RunInput) error { hasYAML := strings.TrimSpace(in.PipelineYAML) != "" hasPath := strings.TrimSpace(in.PipelineYAMLPath) != "" if hasYAML == hasPath { return fmt.Errorf("run request must set exactly one of pipeline_yaml or pipeline_yaml_path") } return nil } // runRequestFromJSON converts an incoming otopb.RunRequest decoded from the job // create body into a store-level RunInput. It returns nil when no run request // was supplied, keeping job creation backward compatible. func runRequestFromJSON(rr *otopb.RunRequest) *cicdstate.RunInput { if rr == nil { return nil } in := &cicdstate.RunInput{ PipelineYAMLPath: rr.GetPipelineYamlPath(), PipelineYAML: rr.GetPipelineYaml(), CommandTypes: append([]string(nil), rr.GetCommandTypes()...), } if len(rr.GetVariables()) > 0 { in.Variables = make(map[string]string, len(rr.GetVariables())) for k, v := range rr.GetVariables() { in.Variables[k] = v } } return in } // runRequestToProto builds the claim response RunRequest from the stored input, // filling in the runner/job/execution identifiers resolved at claim time. func runRequestToProto(in *cicdstate.RunInput, runnerID, jobID, execID string) *otopb.RunRequest { if in == nil { return nil } rr := &otopb.RunRequest{ RunnerId: runnerID, JobId: jobID, ExecutionId: execID, PipelineYamlPath: in.PipelineYAMLPath, PipelineYaml: in.PipelineYAML, CommandTypes: append([]string(nil), in.CommandTypes...), } if len(in.Variables) > 0 { rr.Variables = make(map[string]string, len(in.Variables)) for k, v := range in.Variables { rr.Variables[k] = v } } return rr } func runInputToJSON(in *cicdstate.RunInput) map[string]interface{} { return map[string]interface{}{ "pipeline_yaml_path": in.PipelineYAMLPath, "pipeline_yaml": in.PipelineYAML, "variables": in.Variables, "command_types": in.CommandTypes, } } type execResponse struct { ID string `json:"id"` JobID string `json:"job_id"` State string `json:"state"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` Logs []cicdstate.LogEntry `json:"logs,omitempty"` Artifacts []cicdstate.ArtifactEntry `json:"artifacts,omitempty"` } func execToJSON(e *cicdstate.Execution) map[string]interface{} { return map[string]interface{}{ "id": e.ID, "job_id": e.JobID, "state": e.State, "created_at": e.CreatedAt.Format(time.RFC3339), "updated_at": e.UpdatedAt.Format(time.RFC3339), "execution_id": e.ID, } } type logEntryResponse struct { Timestamp string `json:"timestamp"` Line string `json:"line"` } func logsToJSON(logs []cicdstate.LogEntry) map[string]interface{} { return map[string]interface{}{ "logs": logs, } } type artifactResponse struct { Name string `json:"name"` Path string `json:"path"` } func artifactsToJSON(artifacts []cicdstate.ArtifactEntry) map[string]interface{} { return map[string]interface{}{ "artifacts": artifacts, } } func writeResponse(w http.ResponseWriter, status int, data interface{}) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(status) _ = json.NewEncoder(w).Encode(data) } func handleRunnerCancelExecution(store *cicdstate.Store, registry *runnerregistry.Registry, runnerID string, execID string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req otopb.CancelRunRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.CancelRunResponse{ Success: false, ErrorMessage: "invalid cancel request", }) return } if err := normalizeRunnerExecutionRequest(registry, runnerID, execID, req.GetRunnerId(), req.GetExecutionId()); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.CancelRunResponse{ Success: false, ErrorMessage: err.Error(), }) return } exec, err := store.GetExecution(execID) if err != nil { writeResponse(w, http.StatusNotFound, &otopb.CancelRunResponse{ Success: false, ErrorMessage: "execution not found", }) return } if exec.RunnerID == "" || exec.RunnerID != runnerID { writeResponse(w, http.StatusBadRequest, &otopb.CancelRunResponse{ Success: false, ErrorMessage: "execution owner mismatch", }) return } if err := store.CancelJobExecution(exec.JobID, execID); err != nil { writeResponse(w, http.StatusBadRequest, &otopb.CancelRunResponse{ Success: false, ErrorMessage: err.Error(), }) return } writeResponse(w, http.StatusOK, &otopb.CancelRunResponse{ Success: true, }) } } func handleRunnerStatus(store *cicdstate.Store, registry *runnerregistry.Registry, runnerID string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { record, ok := registry.Snapshot(runnerID) if !ok { writeResponse(w, http.StatusNotFound, map[string]interface{}{ "accepted": false, "message": "runner not found", }) return } // Find currently running or last execution var lastExec *cicdstate.Execution for _, exec := range store.SnapshotExecutions() { if exec.RunnerID == runnerID { if lastExec == nil || exec.CreatedAt.After(lastExec.CreatedAt) { copyExec := exec lastExec = ©Exec } } } currentJobID := "" currentExecID := "" if lastExec != nil { currentJobID = lastExec.JobID currentExecID = lastExec.ID } writeResponse(w, http.StatusOK, map[string]interface{}{ "accepted": true, "runner_id": runnerID, "status": record.Status, "current_job_id": currentJobID, "current_execution_id": currentExecID, "message": record.FailureReason, }) } } func handleRunnerSelfUpdate(store *cicdstate.Store, registry *runnerregistry.Registry, runnerID string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var req otopb.SelfUpdateRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeResponse(w, http.StatusBadRequest, map[string]interface{}{ "success": false, "accepted": false, "deferred": false, "error_message": "invalid self-update request", "message": "invalid self-update request", }) return } // Normalize runner ID if req.RunnerId == "" { req.RunnerId = runnerID } else if req.RunnerId != runnerID { writeResponse(w, http.StatusBadRequest, map[string]interface{}{ "success": false, "accepted": false, "deferred": false, "error_message": "runner id mismatch between path and body", "message": "runner id mismatch between path and body", }) return } if err := ensureRunnerKnown(registry, runnerID); err != nil { writeResponse(w, http.StatusNotFound, map[string]interface{}{ "success": false, "accepted": false, "deferred": false, "error_message": err.Error(), "message": err.Error(), }) return } if !strings.HasPrefix(req.GetDownloadUrl(), "https://") { writeResponse(w, http.StatusBadRequest, map[string]interface{}{ "success": false, "accepted": false, "deferred": false, "error_message": "download_url must use HTTPS scheme", "message": "download_url must use HTTPS scheme", }) return } // Check active executions hasActiveExecution := false for _, exec := range store.SnapshotExecutions() { if exec.RunnerID == runnerID && exec.State == cicdstate.StateRunning { hasActiveExecution = true break } } if hasActiveExecution { writeResponse(w, http.StatusOK, map[string]interface{}{ "success": false, "accepted": false, "deferred": true, "target_version": req.GetVersion(), "message": "deferred: runner is currently running a job", }) return } writeResponse(w, http.StatusOK, map[string]interface{}{ "success": true, "accepted": true, "deferred": false, "target_version": req.GetVersion(), "message": "self-update request accepted", }) } }