package scheduler import ( "encoding/json" "errors" "net/http" "strconv" "aigateway.local/core/internal/identity" "aigateway.local/core/internal/platform/apiresponse" ) type AdminHTTPHandler struct { service *Service identity *identity.Service mux *http.ServeMux } func NewAdminHTTPHandler(service *Service, identityService *identity.Service) *AdminHTTPHandler { h := &AdminHTTPHandler{service: service, identity: identityService, mux: http.NewServeMux()} h.mux.HandleFunc("GET /api/v1/admin/scheduled-tasks", h.list) h.mux.HandleFunc("POST /api/v1/admin/scheduled-tasks", h.create) h.mux.HandleFunc("GET /api/v1/admin/scheduled-tasks/{id}", h.get) h.mux.HandleFunc("PUT /api/v1/admin/scheduled-tasks/{id}", h.update) h.mux.HandleFunc("DELETE /api/v1/admin/scheduled-tasks/{id}", h.delete) h.mux.HandleFunc("POST /api/v1/admin/scheduled-tasks/{id}/start", h.start) h.mux.HandleFunc("POST /api/v1/admin/scheduled-tasks/{id}/pause", h.pause) h.mux.HandleFunc("POST /api/v1/admin/scheduled-tasks/{id}/run", h.runNow) h.mux.HandleFunc("GET /api/v1/admin/scheduled-tasks/{id}/runs", h.runs) return h } func (h *AdminHTTPHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { h.mux.ServeHTTP(w, r) } func (h *AdminHTTPHandler) require(w http.ResponseWriter, r *http.Request, permission string) (identity.Account, bool) { account, err := h.identity.Authenticate(r.Context(), identity.KindAdmin, r.Header.Get("Authorization")) if err != nil { apiresponse.Error(w, http.StatusUnauthorized, "登录状态无效") return identity.Account{}, false } if !identity.HasPermission(account, permission) { apiresponse.Error(w, http.StatusForbidden, "缺少定时任务权限") return identity.Account{}, false } return account, true } func decodeTask(w http.ResponseWriter, r *http.Request, target any) bool { decoder := json.NewDecoder(http.MaxBytesReader(w, r.Body, 2<<20)) decoder.DisallowUnknownFields() if err := decoder.Decode(target); err != nil { apiresponse.Error(w, http.StatusBadRequest, "请求格式无效") return false } return true } func taskError(w http.ResponseWriter, err error) { if errors.Is(err, ErrNotFound) { apiresponse.Error(w, http.StatusNotFound, "定时任务不存在") return } apiresponse.Error(w, http.StatusBadRequest, err.Error()) } func (h *AdminHTTPHandler) list(w http.ResponseWriter, r *http.Request) { if _, ok := h.require(w, r, identity.PermissionScheduledTaskRead); !ok { return } items, err := h.service.List(r.Context()) if err != nil { taskError(w, err) return } apiresponse.OK(w, items) } func (h *AdminHTTPHandler) get(w http.ResponseWriter, r *http.Request) { if _, ok := h.require(w, r, identity.PermissionScheduledTaskRead); !ok { return } item, err := h.service.Get(r.Context(), r.PathValue("id")) if err != nil { taskError(w, err) return } apiresponse.OK(w, item) } func (h *AdminHTTPHandler) create(w http.ResponseWriter, r *http.Request) { account, ok := h.require(w, r, identity.PermissionScheduledTaskManage) if !ok { return } var input TaskInput if !decodeTask(w, r, &input) { return } item, err := h.service.Save(r.Context(), "", input, account.ID) if err != nil { taskError(w, err) return } apiresponse.OK(w, item) } func (h *AdminHTTPHandler) update(w http.ResponseWriter, r *http.Request) { account, ok := h.require(w, r, identity.PermissionScheduledTaskManage) if !ok { return } var input TaskInput if !decodeTask(w, r, &input) { return } item, err := h.service.Save(r.Context(), r.PathValue("id"), input, account.ID) if err != nil { taskError(w, err) return } apiresponse.OK(w, item) } func (h *AdminHTTPHandler) delete(w http.ResponseWriter, r *http.Request) { if _, ok := h.require(w, r, identity.PermissionScheduledTaskManage); !ok { return } if err := h.service.Delete(r.Context(), r.PathValue("id")); err != nil { taskError(w, err) return } apiresponse.OK(w, map[string]bool{"deleted": true}) } func (h *AdminHTTPHandler) setEnabled(w http.ResponseWriter, r *http.Request, enabled bool) { if _, ok := h.require(w, r, identity.PermissionScheduledTaskManage); !ok { return } item, err := h.service.SetEnabled(r.Context(), r.PathValue("id"), enabled) if err != nil { taskError(w, err) return } apiresponse.OK(w, item) } func (h *AdminHTTPHandler) start(w http.ResponseWriter, r *http.Request) { h.setEnabled(w, r, true) } func (h *AdminHTTPHandler) pause(w http.ResponseWriter, r *http.Request) { h.setEnabled(w, r, false) } func (h *AdminHTTPHandler) runNow(w http.ResponseWriter, r *http.Request) { if _, ok := h.require(w, r, identity.PermissionScheduledTaskManage); !ok { return } run, err := h.service.QueueManual(r.Context(), r.PathValue("id")) if err != nil { taskError(w, err) return } apiresponse.OK(w, run) } func (h *AdminHTTPHandler) runs(w http.ResponseWriter, r *http.Request) { if _, ok := h.require(w, r, identity.PermissionScheduledTaskRead); !ok { return } limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) items, err := h.service.Runs(r.Context(), r.PathValue("id"), limit) if err != nil { taskError(w, err) return } apiresponse.OK(w, items) }