HH-846 / HH-799: add task recovery and audit (#27)

This commit is contained in:
2026-08-31 13:19:54 +08:00
parent 0fbfb335c4
commit b117cc738f
20 changed files with 1545 additions and 61 deletions
+11
View File
@@ -385,6 +385,7 @@ func registerHubWithNetwork(app *fiber.App, store hubStore, probe networkExitPro
}
}
app.Get("/api/browsers", serialized(listBrowsers(store, probe, resolve)))
app.Get("/api/browsers/:alias", getBrowser(store))
app.Post("/api/browsers", serialized(createBrowser(store, probe, resolve)))
app.Post("/api/browsers/:alias/:action", serialized(browserAction(store, probe, resolve)))
app.Delete("/api/browsers/:alias", serialized(deleteBrowser(store)))
@@ -477,6 +478,16 @@ func registerHubWithNetwork(app *fiber.App, store hubStore, probe networkExitPro
})
}
func getBrowser(store hubStore) fiber.Handler {
return func(c fiber.Ctx) error {
environment, err := store.GetEnvironmentContext(c.Context(), c.Params("alias"))
if err != nil {
return hubError(c, err)
}
return c.JSON(environment)
}
}
func listNetworkExits(store hubStore) fiber.Handler {
return func(c fiber.Ctx) error {
exits, err := store.ListNetworkExits(c.Context())
+83 -3
View File
@@ -5,6 +5,8 @@ import (
"encoding/json"
"errors"
"io"
"strconv"
"time"
"git.ipao.vip/rogee/creator-hub/internal/hub"
"git.ipao.vip/rogee/creator-hub/internal/phasea"
@@ -38,6 +40,10 @@ type taskRequest struct {
ConfirmationID string `json:"confirmation_id"`
}
type taskVerificationRequest struct {
Result string `json:"result"`
}
func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeStopStore) {
app.Post("/api/phase-a/accounts", func(c fiber.Ctx) error {
var input accountRequest
@@ -193,7 +199,7 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
})
app.Get("/api/phase-a/tasks", func(c fiber.Ctx) error {
tasks, err := store.ListTasks(c.Context(), c.Query("account_id"), c.Query("draft_id"))
tasks, err := store.ListTasksFiltered(c.Context(), c.Query("account_id"), c.Query("draft_id"), c.Query("state"))
if err != nil {
return phaseAError(c, err)
}
@@ -201,13 +207,21 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
})
app.Get("/api/phase-a/tasks/:id", func(c fiber.Ctx) error {
task, err := store.GetTask(c.Context(), c.Params("id"))
task, err := store.GetTaskDetail(c.Context(), c.Params("id"))
if err != nil {
return phaseAError(c, err)
}
return c.JSON(task)
})
app.Get("/api/phase-a/attempts/:id", func(c fiber.Ctx) error {
attempt, err := store.GetTaskAttemptDetail(c.Context(), c.Params("id"))
if err != nil {
return phaseAError(c, err)
}
return c.JSON(attempt)
})
app.Post("/api/phase-a/tasks/:id/cancel", func(c fiber.Ctx) error {
if err := store.CancelTask(c.Context(), c.Params("id")); err != nil {
return phaseAError(c, err)
@@ -215,6 +229,31 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
return c.SendStatus(fiber.StatusNoContent)
})
app.Post("/api/phase-a/tasks/:id/verify", func(c fiber.Ctx) error {
var input taskVerificationRequest
if err := decodePhaseA(c, &input); err != nil {
return phaseAError(c, err)
}
if err := store.VerifyTask(c.Context(), c.Params("id"), input.Result); err != nil {
return phaseAError(c, err)
}
return c.SendStatus(fiber.StatusNoContent)
})
app.Post("/api/phase-a/tasks/:id/resume", func(c fiber.Ctx) error {
if err := store.ResumeTask(c.Context(), c.Params("id")); err != nil {
return phaseAError(c, err)
}
return c.SendStatus(fiber.StatusNoContent)
})
app.Post("/api/phase-a/tasks/:id/finish", func(c fiber.Ctx) error {
if err := store.FinishTask(c.Context(), c.Params("id")); err != nil {
return phaseAError(c, err)
}
return c.SendStatus(fiber.StatusNoContent)
})
app.Post("/api/phase-a/mock/execute", func(c fiber.Ctx) error {
var input struct {
WorkerID string `json:"worker_id"`
@@ -231,7 +270,11 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
})
app.Get("/api/phase-a/audit", func(c fiber.Ctx) error {
events, err := store.Audit(c.Context())
filter, err := auditFilter(c)
if err != nil {
return phaseAError(c, err)
}
events, err := store.ListAudit(c.Context(), filter)
if err != nil {
return phaseAError(c, err)
}
@@ -239,6 +282,43 @@ func registerPhaseA(app *fiber.App, store *phasea.Store, runtimeStore runtimeSto
})
}
func auditFilter(c fiber.Ctx) (phasea.AuditFilter, error) {
filter := phasea.AuditFilter{
AccountID: c.Query("account_id"), TaskID: c.Query("task_id"), AttemptID: c.Query("attempt_id"),
BrowserEnvAlias: c.Query("browser_env_alias"), NetworkExitID: c.Query("network_exit_id"),
EventType: c.Query("event_type"), Page: 1, PageSize: 25,
}
for _, field := range []struct {
value string
target *int
}{{c.Query("page"), &filter.Page}, {c.Query("page_size"), &filter.PageSize}} {
value, target := field.value, field.target
if value == "" {
continue
}
parsed, err := strconv.Atoi(value)
if err != nil {
return phasea.AuditFilter{}, phasea.ErrInvalid
}
*target = parsed
}
for _, field := range []struct {
value string
target **time.Time
}{{c.Query("from"), &filter.From}, {c.Query("to"), &filter.To}} {
value, target := field.value, field.target
if value == "" {
continue
}
parsed, err := time.Parse(time.RFC3339, value)
if err != nil {
return phasea.AuditFilter{}, phasea.ErrInvalid
}
*target = &parsed
}
return filter, nil
}
func resumeBlockReason(account phasea.Account, environment hub.EnvironmentContext, bindingFound bool) string {
switch {
case account.AuthorizationStatus != "authorized":