asynqmon/task_handlers.go

778 lines
23 KiB
Go
Raw Normal View History

package asynqmon
2020-11-24 22:54:00 +08:00
import (
"encoding/json"
"errors"
"log"
2020-11-24 22:54:00 +08:00
"net/http"
"strconv"
"strings"
2021-01-24 04:06:50 +08:00
"time"
2020-11-24 22:54:00 +08:00
"github.com/gorilla/mux"
2021-09-18 22:23:42 +08:00
2021-05-29 05:40:09 +08:00
"github.com/hibiken/asynq"
2020-11-24 22:54:00 +08:00
)
2020-12-02 23:19:06 +08:00
// ****************************************************************************
// This file defines:
// - http.Handler(s) for task related endpoints
// ****************************************************************************
2021-10-04 23:18:00 +08:00
type listActiveTasksResponse struct {
Tasks []*activeTask `json:"tasks"`
Stats *queueStateSnapshot `json:"stats"`
2021-01-24 04:06:50 +08:00
}
func newListActiveTasksHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter) http.HandlerFunc {
2020-11-24 22:54:00 +08:00
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname := vars["qname"]
pageSize, pageNum := getPageOptions(r)
2021-01-24 04:06:50 +08:00
2020-11-24 22:54:00 +08:00
tasks, err := inspector.ListActiveTasks(
2021-05-29 05:40:09 +08:00
qname, asynq.PageSize(pageSize), asynq.Page(pageNum))
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-05-29 05:40:09 +08:00
qinfo, err := inspector.GetQueueInfo(qname)
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-01-24 04:06:50 +08:00
servers, err := inspector.Servers()
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-10-04 23:18:00 +08:00
// m maps taskID to workerInfo.
2021-05-29 05:40:09 +08:00
m := make(map[string]*asynq.WorkerInfo)
2021-01-24 04:06:50 +08:00
for _, srv := range servers {
for _, w := range srv.ActiveWorkers {
2021-05-29 05:40:09 +08:00
if w.Queue == qname {
m[w.TaskID] = w
2021-01-24 04:06:50 +08:00
}
}
}
activeTasks := toActiveTasks(tasks, pf)
2021-01-24 04:06:50 +08:00
for _, t := range activeTasks {
2021-01-28 08:16:38 +08:00
workerInfo, ok := m[t.ID]
2021-01-24 04:06:50 +08:00
if ok {
2021-01-28 08:16:38 +08:00
t.Started = workerInfo.Started.Format(time.RFC3339)
t.Deadline = workerInfo.Deadline.Format(time.RFC3339)
2021-01-24 04:06:50 +08:00
} else {
t.Started = "-"
2021-01-28 08:16:38 +08:00
t.Deadline = "-"
2021-01-24 04:06:50 +08:00
}
}
2021-10-04 23:18:00 +08:00
resp := listActiveTasksResponse{
2021-01-24 04:06:50 +08:00
Tasks: activeTasks,
2021-10-01 02:49:41 +08:00
Stats: toQueueStateSnapshot(qinfo),
2020-11-24 22:54:00 +08:00
}
2021-01-24 04:06:50 +08:00
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2020-11-24 22:54:00 +08:00
}
}
2021-05-29 05:40:09 +08:00
func newCancelActiveTaskHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
2020-12-05 22:47:35 +08:00
return func(w http.ResponseWriter, r *http.Request) {
id := mux.Vars(r)["task_id"]
2021-05-29 05:40:09 +08:00
if err := inspector.CancelProcessing(id); err != nil {
2020-12-05 22:47:35 +08:00
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newCancelAllActiveTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
const batchSize = 100
page := 1
qname := mux.Vars(r)["qname"]
for {
2021-05-29 05:40:09 +08:00
tasks, err := inspector.ListActiveTasks(qname, asynq.Page(page), asynq.PageSize(batchSize))
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
for _, t := range tasks {
2021-06-19 21:21:07 +08:00
if err := inspector.CancelProcessing(t.ID); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
if len(tasks) < batchSize {
break
}
page++
}
w.WriteHeader(http.StatusNoContent)
}
}
type batchCancelTasksRequest struct {
TaskIDs []string `json:"task_ids"`
}
type batchCancelTasksResponse struct {
CanceledIDs []string `json:"canceled_ids"`
ErrorIDs []string `json:"error_ids"`
}
2021-05-29 05:40:09 +08:00
func newBatchCancelActiveTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
r.Body = http.MaxBytesReader(w, r.Body, maxRequestBodySize)
dec := json.NewDecoder(r.Body)
dec.DisallowUnknownFields()
var req batchCancelTasksRequest
if err := dec.Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
resp := batchCancelTasksResponse{
// avoid null in the json response
CanceledIDs: make([]string, 0),
ErrorIDs: make([]string, 0),
}
for _, id := range req.TaskIDs {
2021-05-29 05:40:09 +08:00
if err := inspector.CancelProcessing(id); err != nil {
log.Printf("error: could not send cancelation signal to task %s", id)
resp.ErrorIDs = append(resp.ErrorIDs, id)
} else {
resp.CanceledIDs = append(resp.CanceledIDs, id)
}
}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
func newListPendingTasksHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter) http.HandlerFunc {
2020-11-24 22:54:00 +08:00
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname := vars["qname"]
pageSize, pageNum := getPageOptions(r)
tasks, err := inspector.ListPendingTasks(
2021-05-29 05:40:09 +08:00
qname, asynq.PageSize(pageSize), asynq.Page(pageNum))
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-05-29 05:40:09 +08:00
qinfo, err := inspector.GetQueueInfo(qname)
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
payload := make(map[string]interface{})
if len(tasks) == 0 {
// avoid nil for the tasks field in json output.
2021-10-04 23:18:00 +08:00
payload["tasks"] = make([]*pendingTask, 0)
2020-11-24 22:54:00 +08:00
} else {
payload["tasks"] = toPendingTasks(tasks, pf)
2020-11-24 22:54:00 +08:00
}
2021-10-01 02:49:41 +08:00
payload["stats"] = toQueueStateSnapshot(qinfo)
if err := json.NewEncoder(w).Encode(payload); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2020-11-24 22:54:00 +08:00
}
}
func newListScheduledTasksHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter) http.HandlerFunc {
2020-11-24 22:54:00 +08:00
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname := vars["qname"]
pageSize, pageNum := getPageOptions(r)
tasks, err := inspector.ListScheduledTasks(
2021-05-29 05:40:09 +08:00
qname, asynq.PageSize(pageSize), asynq.Page(pageNum))
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-05-29 05:40:09 +08:00
qinfo, err := inspector.GetQueueInfo(qname)
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
payload := make(map[string]interface{})
if len(tasks) == 0 {
// avoid nil for the tasks field in json output.
2021-10-04 23:18:00 +08:00
payload["tasks"] = make([]*scheduledTask, 0)
2020-11-24 22:54:00 +08:00
} else {
payload["tasks"] = toScheduledTasks(tasks, pf)
2020-11-24 22:54:00 +08:00
}
2021-10-01 02:49:41 +08:00
payload["stats"] = toQueueStateSnapshot(qinfo)
if err := json.NewEncoder(w).Encode(payload); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2020-11-24 22:54:00 +08:00
}
}
func newListRetryTasksHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter) http.HandlerFunc {
2020-11-24 22:54:00 +08:00
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname := vars["qname"]
pageSize, pageNum := getPageOptions(r)
tasks, err := inspector.ListRetryTasks(
2021-05-29 05:40:09 +08:00
qname, asynq.PageSize(pageSize), asynq.Page(pageNum))
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-05-29 05:40:09 +08:00
qinfo, err := inspector.GetQueueInfo(qname)
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
payload := make(map[string]interface{})
if len(tasks) == 0 {
// avoid nil for the tasks field in json output.
2021-10-04 23:18:00 +08:00
payload["tasks"] = make([]*retryTask, 0)
2020-11-24 22:54:00 +08:00
} else {
payload["tasks"] = toRetryTasks(tasks, pf)
2020-11-24 22:54:00 +08:00
}
2021-10-01 02:49:41 +08:00
payload["stats"] = toQueueStateSnapshot(qinfo)
if err := json.NewEncoder(w).Encode(payload); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2020-11-24 22:54:00 +08:00
}
}
func newListArchivedTasksHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter) http.HandlerFunc {
2020-11-24 22:54:00 +08:00
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname := vars["qname"]
pageSize, pageNum := getPageOptions(r)
tasks, err := inspector.ListArchivedTasks(
2021-05-29 05:40:09 +08:00
qname, asynq.PageSize(pageSize), asynq.Page(pageNum))
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-05-29 05:40:09 +08:00
qinfo, err := inspector.GetQueueInfo(qname)
2020-11-24 22:54:00 +08:00
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
payload := make(map[string]interface{})
if len(tasks) == 0 {
// avoid nil for the tasks field in json output.
2021-10-04 23:18:00 +08:00
payload["tasks"] = make([]*archivedTask, 0)
2020-11-24 22:54:00 +08:00
} else {
payload["tasks"] = toArchivedTasks(tasks, pf)
2020-11-24 22:54:00 +08:00
}
2021-10-01 02:49:41 +08:00
payload["stats"] = toQueueStateSnapshot(qinfo)
if err := json.NewEncoder(w).Encode(payload); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2020-11-24 22:54:00 +08:00
}
}
2021-11-07 06:23:10 +08:00
func newListCompletedTasksHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter, rf ResultFormatter) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname := vars["qname"]
pageSize, pageNum := getPageOptions(r)
tasks, err := inspector.ListCompletedTasks(qname, asynq.PageSize(pageSize), asynq.Page(pageNum))
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
qinfo, err := inspector.GetQueueInfo(qname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
payload := make(map[string]interface{})
if len(tasks) == 0 {
// avoid nil for the tasks field in json output.
payload["tasks"] = make([]*completedTask, 0)
} else {
payload["tasks"] = toCompletedTasks(tasks, pf, rf)
}
payload["stats"] = toQueueStateSnapshot(qinfo)
if err := json.NewEncoder(w).Encode(payload); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
func newListAggregatingTasksHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname := vars["qname"]
gname := vars["gname"]
pageSize, pageNum := getPageOptions(r)
tasks, err := inspector.ListAggregatingTasks(
qname, gname, asynq.PageSize(pageSize), asynq.Page(pageNum))
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
qinfo, err := inspector.GetQueueInfo(qname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
payload := make(map[string]interface{})
if len(tasks) == 0 {
// avoid nil for the tasks field in json output.
payload["tasks"] = make([]*aggregatingTask, 0)
} else {
payload["tasks"] = toAggregatingTasks(tasks, pf)
}
payload["stats"] = toQueueStateSnapshot(qinfo)
if err := json.NewEncoder(w).Encode(payload); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2021-05-29 05:40:09 +08:00
func newDeleteTaskHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
2021-05-29 05:40:09 +08:00
qname, taskid := vars["qname"], vars["task_id"]
if qname == "" || taskid == "" {
http.Error(w, "route parameters should not be empty", http.StatusBadRequest)
return
}
2021-05-29 05:40:09 +08:00
if err := inspector.DeleteTask(qname, taskid); err != nil {
// TODO: Handle task not found error and return 404
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newRunTaskHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
2021-05-29 05:40:09 +08:00
qname, taskid := vars["qname"], vars["task_id"]
if qname == "" || taskid == "" {
http.Error(w, "route parameters should not be empty", http.StatusBadRequest)
return
}
2021-05-29 05:40:09 +08:00
if err := inspector.RunTask(qname, taskid); err != nil {
// TODO: Handle task not found error and return 404
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newArchiveTaskHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
2021-05-29 05:40:09 +08:00
qname, taskid := vars["qname"], vars["task_id"]
if qname == "" || taskid == "" {
http.Error(w, "route parameters should not be empty", http.StatusBadRequest)
return
}
2021-05-29 05:40:09 +08:00
if err := inspector.ArchiveTask(qname, taskid); err != nil {
// TODO: Handle task not found error and return 404
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-10-04 23:18:00 +08:00
type deleteAllTasksResponse struct {
// Number of tasks deleted.
Deleted int `json:"deleted"`
}
2021-05-29 05:40:09 +08:00
func newDeleteAllPendingTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
n, err := inspector.DeleteAllPendingTasks(qname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-10-04 23:18:00 +08:00
resp := deleteAllTasksResponse{n}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
func newDeleteAllAggregatingTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname, gname := vars["qname"], vars["gname"]
n, err := inspector.DeleteAllAggregatingTasks(qname, gname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
resp := deleteAllTasksResponse{n}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2021-05-29 05:40:09 +08:00
func newDeleteAllScheduledTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
n, err := inspector.DeleteAllScheduledTasks(qname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-10-04 23:18:00 +08:00
resp := deleteAllTasksResponse{n}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2021-05-29 05:40:09 +08:00
func newDeleteAllRetryTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
n, err := inspector.DeleteAllRetryTasks(qname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-10-04 23:18:00 +08:00
resp := deleteAllTasksResponse{n}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2021-05-29 05:40:09 +08:00
func newDeleteAllArchivedTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
n, err := inspector.DeleteAllArchivedTasks(qname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
2021-10-04 23:18:00 +08:00
resp := deleteAllTasksResponse{n}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2021-11-07 06:23:10 +08:00
func newDeleteAllCompletedTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
n, err := inspector.DeleteAllCompletedTasks(qname)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
resp := deleteAllTasksResponse{n}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2021-05-29 05:40:09 +08:00
func newRunAllScheduledTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
if _, err := inspector.RunAllScheduledTasks(qname); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newRunAllRetryTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
if _, err := inspector.RunAllRetryTasks(qname); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newRunAllArchivedTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
2020-12-16 23:35:36 +08:00
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
if _, err := inspector.RunAllArchivedTasks(qname); err != nil {
2020-12-16 23:35:36 +08:00
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
func newRunAllAggregatingTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname, gname := vars["qname"], vars["gname"]
if _, err := inspector.RunAllAggregatingTasks(qname, gname); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newArchiveAllPendingTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
if _, err := inspector.ArchiveAllPendingTasks(qname); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
func newArchiveAllAggregatingTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname, gname := vars["qname"], vars["gname"]
if _, err := inspector.ArchiveAllAggregatingTasks(qname, gname); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newArchiveAllScheduledTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
if _, err := inspector.ArchiveAllScheduledTasks(qname); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
2021-05-29 05:40:09 +08:00
func newArchiveAllRetryTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
qname := mux.Vars(r)["qname"]
if _, err := inspector.ArchiveAllRetryTasks(qname); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusNoContent)
}
}
// request body used for all batch delete tasks endpoints.
type batchDeleteTasksRequest struct {
2021-05-29 05:40:09 +08:00
TaskIDs []string `json:"task_ids"`
}
// Note: Redis does not have any rollback mechanism, so it's possible
// to have partial success when doing a batch operation.
2021-05-29 05:40:09 +08:00
// For this reason this response contains a list of succeeded ids
// and a list of failed ids.
type batchDeleteTasksResponse struct {
2021-05-29 05:40:09 +08:00
// task ids that were successfully deleted.
DeletedIDs []string `json:"deleted_ids"`
2021-05-29 05:40:09 +08:00
// task ids that were not deleted.
FailedIDs []string `json:"failed_ids"`
}
// Maximum request body size in bytes.
// Allow up to 1MB in size.
const maxRequestBodySize = 1000000
2021-05-29 05:40:09 +08:00
func newBatchDeleteTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
r.Body = http.MaxBytesReader(w, r.Body, maxRequestBodySize)
dec := json.NewDecoder(r.Body)
dec.DisallowUnknownFields()
var req batchDeleteTasksRequest
if err := dec.Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
qname := mux.Vars(r)["qname"]
resp := batchDeleteTasksResponse{
// avoid null in the json response
2021-05-29 05:40:09 +08:00
DeletedIDs: make([]string, 0),
FailedIDs: make([]string, 0),
}
2021-05-29 05:40:09 +08:00
for _, taskid := range req.TaskIDs {
if err := inspector.DeleteTask(qname, taskid); err != nil {
log.Printf("error: could not delete task with id %q: %v", taskid, err)
resp.FailedIDs = append(resp.FailedIDs, taskid)
} else {
2021-05-29 05:40:09 +08:00
resp.DeletedIDs = append(resp.DeletedIDs, taskid)
}
}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2020-12-15 22:16:58 +08:00
type batchRunTasksRequest struct {
2021-05-29 05:40:09 +08:00
TaskIDs []string `json:"task_ids"`
2020-12-15 22:16:58 +08:00
}
type batchRunTasksResponse struct {
2021-05-29 05:40:09 +08:00
// task ids that were successfully moved to the pending state.
PendingIDs []string `json:"pending_ids"`
// task ids that were not able to move to the pending state.
ErrorIDs []string `json:"error_ids"`
2020-12-15 22:16:58 +08:00
}
2021-05-29 05:40:09 +08:00
func newBatchRunTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
2020-12-15 22:16:58 +08:00
return func(w http.ResponseWriter, r *http.Request) {
r.Body = http.MaxBytesReader(w, r.Body, maxRequestBodySize)
dec := json.NewDecoder(r.Body)
dec.DisallowUnknownFields()
var req batchRunTasksRequest
if err := dec.Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
qname := mux.Vars(r)["qname"]
resp := batchRunTasksResponse{
// avoid null in the json response
2021-05-29 05:40:09 +08:00
PendingIDs: make([]string, 0),
ErrorIDs: make([]string, 0),
2020-12-15 22:16:58 +08:00
}
2021-05-29 05:40:09 +08:00
for _, taskid := range req.TaskIDs {
if err := inspector.RunTask(qname, taskid); err != nil {
log.Printf("error: could not run task with id %q: %v", taskid, err)
resp.ErrorIDs = append(resp.ErrorIDs, taskid)
2020-12-15 22:16:58 +08:00
} else {
2021-05-29 05:40:09 +08:00
resp.PendingIDs = append(resp.PendingIDs, taskid)
2020-12-15 22:16:58 +08:00
}
}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
type batchArchiveTasksRequest struct {
2021-05-29 05:40:09 +08:00
TaskIDs []string `json:"task_ids"`
}
type batchArchiveTasksResponse struct {
2021-05-29 05:40:09 +08:00
// task ids that were successfully moved to the archived state.
ArchivedIDs []string `json:"archived_ids"`
// task ids that were not able to move to the archived state.
ErrorIDs []string `json:"error_ids"`
}
2021-05-29 05:40:09 +08:00
func newBatchArchiveTasksHandlerFunc(inspector *asynq.Inspector) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
r.Body = http.MaxBytesReader(w, r.Body, maxRequestBodySize)
dec := json.NewDecoder(r.Body)
dec.DisallowUnknownFields()
var req batchArchiveTasksRequest
if err := dec.Decode(&req); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
qname := mux.Vars(r)["qname"]
resp := batchArchiveTasksResponse{
// avoid null in the json response
2021-05-29 05:40:09 +08:00
ArchivedIDs: make([]string, 0),
ErrorIDs: make([]string, 0),
}
2021-05-29 05:40:09 +08:00
for _, taskid := range req.TaskIDs {
if err := inspector.ArchiveTask(qname, taskid); err != nil {
log.Printf("error: could not archive task with id %q: %v", taskid, err)
resp.ErrorIDs = append(resp.ErrorIDs, taskid)
} else {
2021-05-29 05:40:09 +08:00
resp.ArchivedIDs = append(resp.ArchivedIDs, taskid)
}
}
if err := json.NewEncoder(w).Encode(resp); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}
2020-11-24 22:54:00 +08:00
// getPageOptions read page size and number from the request url if set,
// otherwise it returns the default value.
func getPageOptions(r *http.Request) (pageSize, pageNum int) {
pageSize = 20 // default page size
pageNum = 1 // default page num
q := r.URL.Query()
if s := q.Get("size"); s != "" {
if n, err := strconv.Atoi(s); err == nil {
pageSize = n
}
}
if s := q.Get("page"); s != "" {
if n, err := strconv.Atoi(s); err == nil {
pageNum = n
}
}
return pageSize, pageNum
}
2021-11-07 06:23:10 +08:00
func newGetTaskHandlerFunc(inspector *asynq.Inspector, pf PayloadFormatter, rf ResultFormatter) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
qname, taskid := vars["qname"], vars["task_id"]
if qname == "" {
http.Error(w, "queue name cannot be empty", http.StatusBadRequest)
return
}
if taskid == "" {
http.Error(w, "task_id cannot be empty", http.StatusBadRequest)
return
}
info, err := inspector.GetTaskInfo(qname, taskid)
switch {
case errors.Is(err, asynq.ErrQueueNotFound), errors.Is(err, asynq.ErrTaskNotFound):
http.Error(w, strings.TrimPrefix(err.Error(), "asynq: "), http.StatusNotFound)
return
case err != nil:
http.Error(w, strings.TrimPrefix(err.Error(), "asynq: "), http.StatusInternalServerError)
return
}
2021-11-07 06:23:10 +08:00
if err := json.NewEncoder(w).Encode(toTaskInfo(info, pf, rf)); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
}
}