oto/services/core/internal/httpserver/server.go
toki e418925e61 feat(control-plane): runner 등록 계약을 연결한다
OTO Server가 runner 등록 상태를 직접 소유하고 Dart runner가 Server registration 계약으로 전환되어야 Control Plane 분리 마이그레이션을 이어갈 수 있다.
2026-06-05 15:34:17 +09:00

99 lines
2.6 KiB
Go

package httpserver
import (
"context"
"encoding/json"
"net"
"net/http"
"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 NewServerWithRegistry(addr, runnerregistry.New())
}
// NewServerWithRegistry creates a server using an injected runner registry.
func NewServerWithRegistry(addr string, registry *runnerregistry.Registry) *Server {
mux := http.NewServeMux()
// Register health and readiness endpoints
mux.HandleFunc("/healthz", handleHealthz)
mux.HandleFunc("/readyz", handleReadyz)
mux.HandleFunc("/api/v1/runners/register", handleRunnerRegister(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",
})
return
}
response := registry.Register(&request)
writeRunnerRegisterResponse(w, http.StatusOK, response)
}
}
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)
}