diff --git a/cmd/control-plane/hub.go b/cmd/control-plane/hub.go index 48237ae..33e32ea 100644 --- a/cmd/control-plane/hub.go +++ b/cmd/control-plane/hub.go @@ -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()) diff --git a/cmd/control-plane/phasea.go b/cmd/control-plane/phasea.go index 43670dc..e9aa124 100644 --- a/cmd/control-plane/phasea.go +++ b/cmd/control-plane/phasea.go @@ -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": diff --git a/internal/hub/environment.go b/internal/hub/environment.go index 77e288f..c13d127 100644 --- a/internal/hub/environment.go +++ b/internal/hub/environment.go @@ -326,9 +326,17 @@ func invalidateAccountsForExit(ctx context.Context, tx *sql.Tx, exitID string) e FROM environment_binding binding WHERE binding.network_exit_id = $1 AND binding.account_id = account.id RETURNING account.id + ), held AS ( + UPDATE operation_task task SET + state = CASE task.state WHEN 'executing' THEN 'needs_confirmation' ELSE 'policy_hold' END, + hold_reason = 'exit_unhealthy', verification_result = NULL, verified_at = NULL, verified_by = NULL, + lease_owner = NULL, lease_until = NULL, updated_at = now() + FROM changed WHERE task.account_id = changed.id AND task.state IN ('queued', 'executing') + RETURNING task.current_attempt_id, task.state ) - UPDATE operation_task task SET state = 'policy_hold', updated_at = now() - FROM changed WHERE task.account_id = changed.id AND task.state = 'queued'`, exitID); err != nil { + UPDATE execution_attempt attempt SET finished_at = now(), outcome = 'uncertain' + FROM held WHERE held.state = 'needs_confirmation' AND attempt.id = held.current_attempt_id + AND attempt.finished_at IS NULL`, exitID); err != nil { return errors.New("invalidate network exit accounts") } return nil diff --git a/internal/hub/migration_test.go b/internal/hub/migration_test.go index 9bd30cd..0ce811c 100644 --- a/internal/hub/migration_test.go +++ b/internal/hub/migration_test.go @@ -3,6 +3,7 @@ package hub import ( "context" "database/sql" + "errors" "fmt" "net/url" "os" @@ -29,14 +30,64 @@ func TestUnifiedAccountMigration(t *testing.T) { t.Fatal(err) } defer db.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version IN (1, 2, 3, 4, 5, 6, 7, 8, 9, 10)`, 10) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version BETWEEN 1 AND 12`, 12) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name IN ('social_account', 'browser_env', 'network_exit', 'environment_binding')`, 4) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'environment_binding' AND column_name = 'runtime_cleanup_pending'`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'environment_binding' AND column_name LIKE 'runtime_cleanup_%'`, 5) store = openFullyMigratedHub(t, ctx, testURL) store.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version IN (1, 2, 3, 4, 5, 6, 7, 8, 9, 10)`, 10) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version BETWEEN 1 AND 12`, 12) + }) + + t.Run("previous migration 011 already applied", func(t *testing.T) { + ctx := context.Background() + testURL := isolatedDatabaseURL(t, databaseURL) + store := openFullyMigratedHub(t, ctx, testURL) + store.Close() + db, err := sql.Open("pgx", testURL) + if err != nil { + t.Fatal(err) + } + defer db.Close() + if _, err := db.Exec(` + INSERT INTO credential_reference (id, provider, reference_key) + VALUES ('credential-recovery-upgrade', 'os_keyring', 'creatorhub/recovery-upgrade'); + INSERT INTO social_account + (id, credential_reference_id, platform, platform_account_key, authorization_kind, authorization_status) + VALUES ('recovery-upgrade', 'credential-recovery-upgrade', 'mock', 'recovery-upgrade', 'owned', 'authorized'); + INSERT INTO content_draft (id, account_id, version, content) + VALUES ('draft-recovery-upgrade', 'recovery-upgrade', 1, 'legacy'); + INSERT INTO confirmation (id, account_id, account_version, draft_id, draft_version, version) + VALUES ('confirmation-recovery-upgrade', 'recovery-upgrade', 1, 'draft-recovery-upgrade', 1, 1); + INSERT INTO operation_task + (id, idempotency_key, account_id, account_version, draft_id, draft_version, + confirmation_id, confirmation_version, state, hold_reason) + VALUES ('task-recovery-upgrade', 'task-recovery-upgrade-key', 'recovery-upgrade', 1, + 'draft-recovery-upgrade', 1, 'confirmation-recovery-upgrade', 1, + 'needs_confirmation', 'legacy_confirmation_required'); + DELETE FROM schema_migration WHERE version = 12`); err != nil { + t.Fatal(err) + } + + store, err = Open(ctx, testURL) + if err != nil { + t.Fatalf("apply task recovery compatibility migration: %v", err) + } + store.Close() + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration WHERE version = 12`, 1) + assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task + WHERE id = 'task-recovery-upgrade' AND state = 'needs_confirmation' + AND hold_reason = 'task_result_uncertain'`, 1) + phaseAStore, err := phasea.Open(ctx, testURL) + if err != nil { + t.Fatal(err) + } + defer phaseAStore.Close() + detail, err := phaseAStore.GetTaskDetail(ctx, "task-recovery-upgrade") + if err != nil || detail.AllowedAction != "verify" { + t.Fatalf("migrated unknown result did not require manual verification: detail=%+v err=%v", detail, err) + } }) t.Run("previous migration 008 already applied", func(t *testing.T) { @@ -133,8 +184,10 @@ func TestUnifiedAccountMigration(t *testing.T) { INSERT INTO confirmation (id, account_id, account_version, draft_id, draft_version, version) VALUES ('legacy-confirmation', 'mapped', 1, 'legacy-draft', 1, 1); INSERT INTO operation_task - (id, idempotency_key, account_id, account_version, draft_id, draft_version, confirmation_id, confirmation_version) - VALUES ('legacy-task', 'legacy-task-key', 'mapped', 1, 'legacy-draft', 1, 'legacy-confirmation', 1); + (id, idempotency_key, account_id, account_version, draft_id, draft_version, confirmation_id, confirmation_version, state) + VALUES + ('legacy-task', 'legacy-task-key', 'mapped', 1, 'legacy-draft', 1, 'legacy-confirmation', 1, 'queued'), + ('legacy-unknown-task', 'legacy-unknown-task-key', 'mapped', 1, 'legacy-draft', 1, 'legacy-confirmation', 1, 'needs_confirmation'); INSERT INTO gateway (name, endpoint, token) VALUES ('legacy-gateway', 'http://127.0.0.1:8081', 'legacy-gateway-token'); INSERT INTO browser_image (version, image_ref) VALUES ('1', 'example/browser:1'); INSERT INTO browser_env (alias, name, gateway_name, image_version, fingerprint) VALUES @@ -163,6 +216,7 @@ func TestUnifiedAccountMigration(t *testing.T) { assertDatabaseCount(t, db, `SELECT count(*) FROM runtime_instance WHERE id = 'instance-unbound' AND binding_id IS NULL`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM audit_event WHERE event_type = 'legacy_event'`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task WHERE id = 'legacy-task' AND state = 'policy_hold'`, 1) + assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task WHERE id = 'legacy-unknown-task' AND state = 'needs_confirmation' AND hold_reason = 'task_result_uncertain'`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM browser_env WHERE alias = 'mapped' AND NOT (fingerprint ?| ARRAY['proxy_server', 'disable_non_proxied_udp'])`, 1) if _, err := db.Exec(`UPDATE environment_binding SET runtime_cleanup_pending = true WHERE id = 'mapped'`); err != nil { t.Fatalf("migration 008 blocked an old writer setting cleanup pending: %v", err) @@ -227,7 +281,13 @@ func TestUnifiedAccountMigration(t *testing.T) { VALUES ('upgrade-confirmation', 'mapped', 1, 'upgrade-draft', 1, 1); INSERT INTO operation_task (id, idempotency_key, account_id, account_version, draft_id, draft_version, confirmation_id, confirmation_version) - VALUES ('upgrade-task', 'upgrade-task-key', 'mapped', 1, 'upgrade-draft', 1, 'upgrade-confirmation', 1)`); err != nil { + VALUES ('upgrade-task', 'upgrade-task-key', 'mapped', 1, 'upgrade-draft', 1, 'upgrade-confirmation', 1); + INSERT INTO operation_task + (id, idempotency_key, account_id, account_version, draft_id, draft_version, confirmation_id, confirmation_version, state, lease_owner, lease_until) + VALUES ('upgrade-executing', 'upgrade-executing-key', 'mapped', 1, 'upgrade-draft', 1, 'upgrade-confirmation', 1, + 'executing', 'worker-old', now() + interval '1 minute'); + INSERT INTO execution_attempt (id, task_id) VALUES ('upgrade-attempt', 'upgrade-executing'); + UPDATE operation_task SET current_attempt_id = 'upgrade-attempt' WHERE id = 'upgrade-executing'`); err != nil { t.Fatal(err) } store, err = Open(ctx, testURL) @@ -241,7 +301,26 @@ func TestUnifiedAccountMigration(t *testing.T) { assertDatabaseCount(t, db, `SELECT count(*) FROM browser_env WHERE alias = 'mapped' AND version = 2 AND image_version = '2'`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM environment_binding WHERE id = 'mapped' AND version = 2`, 1) assertDatabaseCount(t, db, `SELECT count(*) FROM social_account WHERE id = 'mapped' AND version = 2 AND status = 'paused'`, 1) - assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task WHERE id = 'upgrade-task' AND state = 'policy_hold'`, 1) + assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task WHERE id = 'upgrade-task' AND state = 'policy_hold' AND hold_reason = 'binding_version_changed'`, 1) + assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task WHERE id = 'upgrade-executing' AND state = 'needs_confirmation' AND hold_reason = 'task_result_uncertain' AND lease_owner IS NULL`, 1) + assertDatabaseCount(t, db, `SELECT count(*) FROM execution_attempt WHERE id = 'upgrade-attempt' AND outcome = 'uncertain' AND finished_at IS NOT NULL`, 1) + phaseAStore, err = phasea.Open(ctx, testURL) + if err != nil { + t.Fatal(err) + } + defer phaseAStore.Close() + detail, err := phaseAStore.GetTaskDetail(ctx, "upgrade-executing") + if err != nil || detail.AllowedAction != "verify" { + t.Fatalf("upgraded executing task did not require manual verification: detail=%+v err=%v", detail, err) + } + if err := phaseAStore.VerifyTask(ctx, "upgrade-executing", "failed"); err != nil { + t.Fatalf("verify upgraded executing task: %v", err) + } + assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task WHERE id = 'upgrade-executing' AND verification_result = 'failed' AND verified_by = 'local-user'`, 1) + if err := phaseAStore.VerifyTask(ctx, "upgrade-executing", "succeeded"); !errors.Is(err, phasea.ErrConflict) { + t.Fatalf("repeated verification changed the recorded conclusion: %v", err) + } + assertDatabaseCount(t, db, `SELECT count(*) FROM operation_task WHERE confirmation_id = 'upgrade-confirmation'`, 2) store, err = Open(ctx, testURL) if err != nil { diff --git a/internal/hub/migrations/011_task_recovery.sql b/internal/hub/migrations/011_task_recovery.sql new file mode 100644 index 0000000..5b16550 --- /dev/null +++ b/internal/hub/migrations/011_task_recovery.sql @@ -0,0 +1,27 @@ +ALTER TABLE operation_task + ADD COLUMN hold_reason text, + ADD COLUMN verification_result text CHECK (verification_result IN ('not_executed', 'succeeded', 'failed')), + ADD COLUMN verified_at timestamptz, + ADD COLUMN verified_by text, + ADD CONSTRAINT operation_task_verification_consistent CHECK ( + (verification_result IS NULL AND verified_at IS NULL AND verified_by IS NULL) + OR + (verification_result IS NOT NULL AND verified_at IS NOT NULL AND verified_by = 'local-user') + ); + +UPDATE operation_task +SET hold_reason = CASE state + WHEN 'policy_hold' THEN 'legacy_policy_hold' + WHEN 'needs_confirmation' THEN 'task_result_uncertain' +END +WHERE state IN ('policy_hold', 'needs_confirmation'); + +ALTER TABLE execution_attempt + DROP CONSTRAINT execution_attempt_task_id_key; + +CREATE INDEX execution_attempt_task_id_idx + ON execution_attempt (task_id, started_at, id); +CREATE INDEX operation_task_filter_idx + ON operation_task (state, account_id, created_at DESC); +CREATE INDEX audit_event_filters_idx + ON audit_event (created_at DESC, account_id, task_id, event_type); diff --git a/internal/hub/migrations/012_task_recovery_compatibility.sql b/internal/hub/migrations/012_task_recovery_compatibility.sql new file mode 100644 index 0000000..30ab554 --- /dev/null +++ b/internal/hub/migrations/012_task_recovery_compatibility.sql @@ -0,0 +1,4 @@ +UPDATE operation_task +SET hold_reason = 'task_result_uncertain', updated_at = now() +WHERE state = 'needs_confirmation' + AND hold_reason = 'legacy_confirmation_required'; diff --git a/internal/hub/store.go b/internal/hub/store.go index 7adc9ea..ecb4e97 100644 --- a/internal/hub/store.go +++ b/internal/hub/store.go @@ -46,6 +46,12 @@ var migration009 string //go:embed migrations/010_runtime_network_generation.sql var migration010 string +//go:embed migrations/011_task_recovery.sql +var migration011 string + +//go:embed migrations/012_task_recovery_compatibility.sql +var migration012 string + var ( ErrConflict = errors.New("resource conflicts with existing state") ErrInvalid = errors.New("invalid hub input") @@ -131,7 +137,7 @@ func (s *Store) migrate(ctx context.Context) error { for _, migration := range []struct { version int sql string - }{{2, migration002}, {3, migration003}, {4, migration004}, {5, migration005}, {6, migration006}, {7, migration007}, {8, migration008}, {9, migration009}, {10, migration010}} { + }{{2, migration002}, {3, migration003}, {4, migration004}, {5, migration005}, {6, migration006}, {7, migration007}, {8, migration008}, {9, migration009}, {10, migration010}, {11, migration011}, {12, migration012}} { var applied bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil { return errors.New("read hub schema migration state") @@ -394,9 +400,18 @@ func (s *Store) UpgradeEnv(ctx context.Context, alias, version string) error { FROM environment_binding binding WHERE binding.browser_env_alias = $1 AND binding.account_id = account.id RETURNING account.id + ), held AS ( + UPDATE operation_task task SET + state = CASE task.state WHEN 'executing' THEN 'needs_confirmation' ELSE 'policy_hold' END, + hold_reason = CASE task.state WHEN 'executing' THEN 'task_result_uncertain' ELSE 'binding_version_changed' END, + verification_result = NULL, verified_at = NULL, verified_by = NULL, + lease_owner = NULL, lease_until = NULL, updated_at = now() + FROM changed WHERE task.account_id = changed.id AND task.state IN ('queued', 'executing') + RETURNING task.current_attempt_id, task.state ) - UPDATE operation_task task SET state = 'policy_hold', updated_at = now() - FROM changed WHERE task.account_id = changed.id AND task.state = 'queued'`, alias); err != nil { + UPDATE execution_attempt attempt SET finished_at = now(), outcome = 'uncertain' + FROM held WHERE held.state = 'needs_confirmation' AND attempt.id = held.current_attempt_id + AND attempt.finished_at IS NULL`, alias); err != nil { return errors.New("invalidate upgraded environment account") } return commitHub(tx) diff --git a/internal/hub/store_test.go b/internal/hub/store_test.go index d003534..4e7e8d6 100644 --- a/internal/hub/store_test.go +++ b/internal/hub/store_test.go @@ -502,6 +502,36 @@ func TestNetworkExitBindingRuntimeAndAuditWorkflow(t *testing.T) { t.Fatalf("invalid image version must not reach audit persistence: %v", err) } assertDatabaseCount(t, store.db, `SELECT count(*) FROM audit_event WHERE operation_id = '`+invalidAction.OperationID+`'`, 0) + var accountVersion int64 + if err := store.db.QueryRowContext(ctx, `SELECT version FROM social_account WHERE id = 'account-a'`).Scan(&accountVersion); err != nil { + t.Fatal(err) + } + if _, err := store.db.ExecContext(ctx, `INSERT INTO content_draft (id, account_id, version, content) VALUES ('exit-hold-draft', 'account-a', 1, 'test')`); err != nil { + t.Fatal(err) + } + if _, err := store.db.ExecContext(ctx, `INSERT INTO confirmation (id, account_id, account_version, draft_id, draft_version, version) + VALUES ('exit-hold-confirmation', 'account-a', $1, 'exit-hold-draft', 1, 1)`, accountVersion); err != nil { + t.Fatal(err) + } + if _, err := store.db.ExecContext(ctx, `INSERT INTO operation_task + (id, idempotency_key, account_id, account_version, draft_id, draft_version, confirmation_id, confirmation_version, state, lease_owner, lease_until) + VALUES + ('exit-hold-queued', 'exit-hold-queued-key', 'account-a', $1, 'exit-hold-draft', 1, 'exit-hold-confirmation', 1, 'queued', NULL, NULL), + ('exit-hold-executing', 'exit-hold-executing-key', 'account-a', $1, 'exit-hold-draft', 1, 'exit-hold-confirmation', 1, 'executing', 'worker-old', now() + interval '1 minute')`, accountVersion); err != nil { + t.Fatal(err) + } + if _, err := store.db.ExecContext(ctx, `INSERT INTO execution_attempt (id, task_id) VALUES ('exit-hold-attempt', 'exit-hold-executing')`); err != nil { + t.Fatal(err) + } + if _, err := store.db.ExecContext(ctx, `UPDATE operation_task SET current_attempt_id = 'exit-hold-attempt' WHERE id = 'exit-hold-executing'`); err != nil { + t.Fatal(err) + } + if _, _, err := store.RecordNetworkExitCheck(ctx, newGeneration.Exit.ID, ExitObservation{}, "proxy_check_failed"); err != nil { + t.Fatal(err) + } + assertDatabaseCount(t, store.db, `SELECT count(*) FROM operation_task WHERE id = 'exit-hold-queued' AND state = 'policy_hold' AND hold_reason = 'exit_unhealthy'`, 1) + assertDatabaseCount(t, store.db, `SELECT count(*) FROM operation_task WHERE id = 'exit-hold-executing' AND state = 'needs_confirmation' AND hold_reason = 'exit_unhealthy' AND lease_owner IS NULL`, 1) + assertDatabaseCount(t, store.db, `SELECT count(*) FROM execution_attempt WHERE id = 'exit-hold-attempt' AND outcome = 'uncertain' AND finished_at IS NOT NULL`, 1) var auditText string if err := store.db.QueryRowContext(ctx, `SELECT string_agg(row_to_json(event)::text, '') FROM audit_event event`).Scan(&auditText); err != nil { t.Fatal(err) diff --git a/internal/phasea/store.go b/internal/phasea/store.go index 90c5d49..106b896 100644 --- a/internal/phasea/store.go +++ b/internal/phasea/store.go @@ -29,6 +29,7 @@ var ( refPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._/-]{0,127}$`) platformKeyPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._:@/-]{0,127}$`) credentialKeyPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,63}/[A-Za-z0-9][A-Za-z0-9._/-]{0,126}$`) + eventPattern = regexp.MustCompile(`^[a-z0-9_]{1,64}$`) ) type Store struct{ db *sql.DB } @@ -69,16 +70,50 @@ type Confirmation struct { } type Task struct { - ID string `json:"id"` - IdempotencyKey string `json:"idempotency_key"` - AccountID string `json:"account_id"` - AccountVersion int64 `json:"account_version"` - DraftID string `json:"draft_id"` - DraftVersion int64 `json:"draft_version"` - ConfirmationID string `json:"confirmation_id"` - ConfirmationVersion int64 `json:"confirmation_version"` - State string `json:"state"` - CreatedAt time.Time `json:"created_at"` + ID string `json:"id"` + IdempotencyKey string `json:"idempotency_key"` + AccountID string `json:"account_id"` + AccountVersion int64 `json:"account_version"` + DraftID string `json:"draft_id"` + DraftVersion int64 `json:"draft_version"` + ConfirmationID string `json:"confirmation_id"` + ConfirmationVersion int64 `json:"confirmation_version"` + State string `json:"state"` + HoldReason string `json:"hold_reason,omitempty"` + VerificationResult string `json:"verification_result,omitempty"` + VerifiedAt *time.Time `json:"verified_at,omitempty"` + VerifiedBy string `json:"verified_by,omitempty"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +type TaskAttempt struct { + ID string `json:"id"` + TaskID string `json:"task_id,omitempty"` + StartedAt time.Time `json:"started_at"` + FinishedAt *time.Time `json:"finished_at,omitempty"` + Outcome string `json:"outcome,omitempty"` + Evidence map[string]string `json:"evidence"` +} + +type TaskAttemptDetail struct { + TaskAttempt + BrowserEnvAlias string `json:"browser_env_alias,omitempty"` + NetworkExitID string `json:"network_exit_id,omitempty"` + RuntimeInstanceID string `json:"runtime_instance_id,omitempty"` + BindingVersion int64 `json:"binding_version,omitempty"` +} + +type TaskDetail struct { + Task + Confirmation ConfirmationSnapshot `json:"confirmation"` + Attempts []TaskAttempt `json:"attempts"` + BrowserEnvAlias string `json:"browser_env_alias,omitempty"` + NetworkExitID string `json:"network_exit_id,omitempty"` + RuntimeInstanceID string `json:"runtime_instance_id,omitempty"` + BindingVersion int64 `json:"binding_version,omitempty"` + AllowedAction string `json:"allowed_action,omitempty"` + ReadinessReason string `json:"readiness_reason,omitempty"` } type ConfirmationSnapshot struct { @@ -137,6 +172,19 @@ type AuditEvent struct { CreatedAt time.Time `json:"created_at"` } +type AuditFilter struct { + AccountID, TaskID, AttemptID, BrowserEnvAlias, NetworkExitID, EventType string + From, To *time.Time + Page, PageSize int +} + +type AuditPage struct { + Data []AuditEvent `json:"data"` + Total int `json:"total"` + Page int `json:"page"` + PageSize int `json:"page_size"` +} + func Open(ctx context.Context, databaseURL string) (*Store, error) { db, err := sql.Open("pgx", databaseURL) if err != nil { @@ -529,14 +577,22 @@ func scanConfirmation(row confirmationScanner) (ConfirmationSnapshot, error) { } func (s *Store) ListTasks(ctx context.Context, accountID, draftID string) ([]Task, error) { - if (accountID != "" && !idPattern.MatchString(accountID)) || (draftID != "" && !refPattern.MatchString(draftID)) { + return s.ListTasksFiltered(ctx, accountID, draftID, "") +} + +func (s *Store) ListTasksFiltered(ctx context.Context, accountID, draftID, state string) ([]Task, error) { + if (accountID != "" && !idPattern.MatchString(accountID)) || (draftID != "" && !refPattern.MatchString(draftID)) || + (state != "" && state != "queued" && state != "executing" && state != "succeeded" && state != "failed" && + state != "needs_confirmation" && state != "policy_hold" && state != "cancelled") { return nil, ErrInvalid } rows, err := s.db.QueryContext(ctx, ` SELECT id, idempotency_key, account_id, account_version, draft_id, draft_version, - confirmation_id, confirmation_version, state, created_at + confirmation_id, confirmation_version, state, hold_reason, verification_result, verified_at, verified_by, + created_at, updated_at FROM operation_task WHERE ($1 = '' OR account_id = $1) AND ($2 = '' OR draft_id = $2) - ORDER BY created_at DESC, id`, accountID, draftID) + AND ($3 = '' OR state = $3) + ORDER BY created_at DESC, id`, accountID, draftID, state) if err != nil { return nil, errors.New("read tasks") } @@ -558,7 +614,8 @@ func (s *Store) GetTask(ctx context.Context, id string) (Task, error) { } return scanTask(s.db.QueryRowContext(ctx, ` SELECT id, idempotency_key, account_id, account_version, draft_id, draft_version, - confirmation_id, confirmation_version, state, created_at + confirmation_id, confirmation_version, state, hold_reason, verification_result, verified_at, verified_by, + created_at, updated_at FROM operation_task WHERE id = $1`, id)) } @@ -566,16 +623,161 @@ type taskScanner interface{ Scan(...any) error } func scanTask(row taskScanner) (Task, error) { var task Task - var confirmationID sql.NullString + var confirmationID, holdReason, verificationResult, verifiedBy sql.NullString var confirmationVersion sql.NullInt64 + var verifiedAt sql.NullTime if err := row.Scan(&task.ID, &task.IdempotencyKey, &task.AccountID, &task.AccountVersion, &task.DraftID, - &task.DraftVersion, &confirmationID, &confirmationVersion, &task.State, &task.CreatedAt); err != nil { + &task.DraftVersion, &confirmationID, &confirmationVersion, &task.State, &holdReason, &verificationResult, + &verifiedAt, &verifiedBy, &task.CreatedAt, &task.UpdatedAt); err != nil { return Task{}, rowError(err) } task.ConfirmationID, task.ConfirmationVersion = confirmationID.String, confirmationVersion.Int64 + task.HoldReason, task.VerificationResult, task.VerifiedBy = holdReason.String, verificationResult.String, verifiedBy.String + if verifiedAt.Valid { + task.VerifiedAt = &verifiedAt.Time + } return task, nil } +func (s *Store) GetTaskDetail(ctx context.Context, id string) (TaskDetail, error) { + task, err := s.GetTask(ctx, id) + if err != nil { + return TaskDetail{}, err + } + var confirmation ConfirmationSnapshot + if task.ConfirmationID != "" { + confirmation, err = s.GetConfirmation(ctx, task.ConfirmationID) + if err != nil { + return TaskDetail{}, err + } + } + rows, err := s.db.QueryContext(ctx, ` + SELECT id, started_at, finished_at, outcome, result + FROM execution_attempt WHERE task_id = $1 ORDER BY started_at, id`, id) + if err != nil { + return TaskDetail{}, errors.New("read task attempts") + } + defer rows.Close() + attempts := []TaskAttempt{} + for rows.Next() { + var attempt TaskAttempt + var finishedAt sql.NullTime + var outcome sql.NullString + var result json.RawMessage + if err := rows.Scan(&attempt.ID, &attempt.StartedAt, &finishedAt, &outcome, &result); err != nil { + return TaskDetail{}, errors.New("decode task attempt") + } + if finishedAt.Valid { + attempt.FinishedAt = &finishedAt.Time + } + attempt.Outcome, attempt.Evidence = outcome.String, safeEvidence(result) + attempts = append(attempts, attempt) + } + if err := rows.Err(); err != nil { + return TaskDetail{}, errors.New("read task attempts") + } + detail := TaskDetail{Task: task, Confirmation: confirmation, Attempts: attempts, AllowedAction: taskAllowedAction(task)} + if task.VerificationResult == "not_executed" { + reason, _, err := taskReadinessReason(ctx, s.db, task.ID) + if err != nil { + return TaskDetail{}, err + } + detail.ReadinessReason = reason + if reason != "" { + detail.AllowedAction = "" + switch reason { + case "account_version_changed", "draft_version_changed", "confirmation_version_changed", "binding_version_changed": + detail.AllowedAction = "reconfirm" + case "confirmation_missing": + if task.ConfirmationID != "" { + detail.AllowedAction = "reconfirm" + } + } + } + } + var browser, network, runtime sql.NullString + var bindingVersion sql.NullInt64 + err = s.db.QueryRowContext(ctx, ` + SELECT browser_env_alias, network_exit_id, runtime_instance_id, binding_version + FROM audit_event WHERE task_id = $1 AND event_type = 'task_claimed' ORDER BY id DESC LIMIT 1`, id). + Scan(&browser, &network, &runtime, &bindingVersion) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return TaskDetail{}, errors.New("read task correlation") + } + detail.BrowserEnvAlias, detail.NetworkExitID = browser.String, network.String + detail.RuntimeInstanceID, detail.BindingVersion = runtime.String, bindingVersion.Int64 + return detail, nil +} + +func (s *Store) GetTaskAttemptDetail(ctx context.Context, id string) (TaskAttemptDetail, error) { + if !refPattern.MatchString(id) { + return TaskAttemptDetail{}, ErrInvalid + } + var detail TaskAttemptDetail + var finishedAt sql.NullTime + var outcome, browser, network, runtime sql.NullString + var bindingVersion sql.NullInt64 + var result json.RawMessage + err := s.db.QueryRowContext(ctx, ` + SELECT attempt.id, attempt.task_id, attempt.started_at, attempt.finished_at, attempt.outcome, attempt.result, + audit.browser_env_alias, audit.network_exit_id, audit.runtime_instance_id, audit.binding_version + FROM execution_attempt attempt + LEFT JOIN LATERAL ( + SELECT browser_env_alias, network_exit_id, runtime_instance_id, binding_version + FROM audit_event WHERE attempt_id = attempt.id AND event_type = 'task_claimed' ORDER BY id DESC LIMIT 1 + ) audit ON true + WHERE attempt.id = $1`, id).Scan(&detail.ID, &detail.TaskID, &detail.StartedAt, &finishedAt, &outcome, &result, + &browser, &network, &runtime, &bindingVersion) + if err != nil { + return TaskAttemptDetail{}, rowError(err) + } + if finishedAt.Valid { + detail.FinishedAt = &finishedAt.Time + } + detail.Outcome, detail.Evidence = outcome.String, safeEvidence(result) + detail.BrowserEnvAlias, detail.NetworkExitID = browser.String, network.String + detail.RuntimeInstanceID, detail.BindingVersion = runtime.String, bindingVersion.Int64 + return detail, nil +} + +func safeEvidence(raw json.RawMessage) map[string]string { + var values map[string]any + if json.Unmarshal(raw, &values) != nil { + return map[string]string{} + } + evidence := map[string]string{} + if outcome, ok := values["mock_outcome"].(string); ok { + evidence["mock_outcome"] = outcome + } + return evidence +} + +func taskAllowedAction(task Task) string { + if task.State != "policy_hold" && task.State != "needs_confirmation" { + return "" + } + if task.VerificationResult == "not_executed" { + return "resume" + } + if task.VerificationResult == "succeeded" || task.VerificationResult == "failed" { + return "finish" + } + switch task.HoldReason { + case "confirmation_missing": + if task.ConfirmationID == "" { + return "" + } + return "reconfirm" + case "account_version_changed", "draft_version_changed", "confirmation_version_changed", "binding_version_changed": + return "reconfirm" + case "account_paused", "account_revoked", "binding_missing", "environment_missing", "exit_missing", "exit_unhealthy", + "runtime_stop_pending", "runtime_missing", "runtime_lease_expired", "execution_lease_expired", "task_result_uncertain", "task_policy_hold": + return "verify" + default: + return "" + } +} + func (s *Store) EnqueueConfirmation(ctx context.Context, confirmationID string) (Task, bool, error) { if !refPattern.MatchString(confirmationID) { return Task{}, false, ErrInvalid @@ -589,7 +791,8 @@ func (s *Store) EnqueueConfirmation(ctx context.Context, confirmationID string) defer tx.Rollback() if existing, err := scanTask(tx.QueryRowContext(ctx, ` SELECT id, idempotency_key, account_id, account_version, draft_id, draft_version, - confirmation_id, confirmation_version, state, created_at + confirmation_id, confirmation_version, state, hold_reason, verification_result, verified_at, verified_by, + created_at, updated_at FROM operation_task WHERE idempotency_key = $1`, idempotencyKey)); err == nil { if err := commit(tx); err != nil { return Task{}, false, err @@ -715,6 +918,7 @@ func enqueueTask(ctx context.Context, tx *sql.Tx, task Task) (Task, bool, error) } if insertedID != "" { task.State = "queued" + task.UpdatedAt = task.CreatedAt if err := appendAudit(ctx, tx, "task_queued", "task_queued", task.AccountID, task.ConfirmationID, task.ConfirmationVersion, "", task.ID, nil); err != nil { return Task{}, false, err } @@ -722,7 +926,8 @@ func enqueueTask(ctx context.Context, tx *sql.Tx, task Task) (Task, bool, error) } existing, err := scanTask(tx.QueryRowContext(ctx, ` SELECT id, idempotency_key, account_id, account_version, draft_id, draft_version, - confirmation_id, confirmation_version, state, created_at + confirmation_id, confirmation_version, state, hold_reason, verification_result, verified_at, verified_by, + created_at, updated_at FROM operation_task WHERE idempotency_key = $1`, task.IdempotencyKey)) if err != nil { return Task{}, false, err @@ -768,14 +973,14 @@ func (s *Store) disableAccount(ctx context.Context, accountID string, revoke boo return errors.New("change account state") } } - held, err := holdQueuedTasks(ctx, tx, accountID) - if err != nil { - return err - } reason := "account_paused" if revoke { reason = "account_revoked" } + held, err := holdQueuedTasks(ctx, tx, accountID, reason) + if err != nil { + return err + } interrupted, err := interruptExecutingTasks(ctx, tx, accountID, reason) if err != nil { return err @@ -844,10 +1049,11 @@ func (s *Store) ResumeAccount(ctx context.Context, accountID string) error { return commit(tx) } -func holdQueuedTasks(ctx context.Context, tx *sql.Tx, accountID string) (int64, error) { +func holdQueuedTasks(ctx context.Context, tx *sql.Tx, accountID, reason string) (int64, error) { result, err := tx.ExecContext(ctx, ` - UPDATE operation_task SET state = 'policy_hold', updated_at = now() - WHERE account_id = $1 AND state = 'queued'`, accountID) + UPDATE operation_task SET state = 'policy_hold', hold_reason = $2, + verification_result = NULL, verified_at = NULL, verified_by = NULL, updated_at = now() + WHERE account_id = $1 AND state = 'queued'`, accountID, reason) if err != nil { return 0, errors.New("hold queued account tasks") } @@ -857,9 +1063,11 @@ func holdQueuedTasks(ctx context.Context, tx *sql.Tx, accountID string) (int64, func interruptExecutingTasks(ctx context.Context, tx *sql.Tx, accountID, reason string) (int64, error) { rows, err := tx.QueryContext(ctx, ` - UPDATE operation_task SET state = 'needs_confirmation', lease_owner = NULL, lease_until = NULL, updated_at = now() + UPDATE operation_task SET state = 'needs_confirmation', hold_reason = $2, + verification_result = NULL, verified_at = NULL, verified_by = NULL, + lease_owner = NULL, lease_until = NULL, updated_at = now() WHERE account_id = $1 AND state = 'executing' - RETURNING id, current_attempt_id, confirmation_id, confirmation_version`, accountID) + RETURNING id, current_attempt_id, confirmation_id, confirmation_version`, accountID, reason) if err != nil { return 0, errors.New("interrupt executing account tasks") } @@ -898,6 +1106,191 @@ func interruptExecutingTasks(ctx context.Context, tx *sql.Tx, accountID, reason return int64(len(tasks)), nil } +func (s *Store) VerifyTask(ctx context.Context, taskID, result string) error { + if !refPattern.MatchString(taskID) || (result != "not_executed" && result != "succeeded" && result != "failed") { + return ErrInvalid + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return errors.New("begin task verification transaction") + } + defer tx.Rollback() + task, err := scanTask(tx.QueryRowContext(ctx, ` + SELECT id, idempotency_key, account_id, account_version, draft_id, draft_version, + confirmation_id, confirmation_version, state, hold_reason, verification_result, verified_at, verified_by, + created_at, updated_at + FROM operation_task WHERE id = $1 FOR UPDATE`, taskID)) + if err != nil { + return err + } + if task.State != "policy_hold" && task.State != "needs_confirmation" { + return ErrConflict + } + if task.VerificationResult != "" { + return ErrConflict + } + switch task.HoldReason { + case "account_paused", "account_revoked", "binding_missing", "environment_missing", "exit_missing", "exit_unhealthy", + "runtime_stop_pending", "runtime_missing", "runtime_lease_expired", "execution_lease_expired", "task_result_uncertain", "task_policy_hold": + default: + return ErrConflict + } + if _, err := tx.ExecContext(ctx, ` + UPDATE operation_task SET verification_result = $2, verified_at = now(), verified_by = 'local-user', updated_at = now() + WHERE id = $1`, taskID, result); err != nil { + return errors.New("record task verification") + } + if err := appendAudit(ctx, tx, "task_verified", "manual_verification_recorded", task.AccountID, task.ConfirmationID, + task.ConfirmationVersion, "", task.ID, map[string]string{"verification_result": result}); err != nil { + return err + } + return commit(tx) +} + +func (s *Store) ResumeTask(ctx context.Context, taskID string) error { + if !refPattern.MatchString(taskID) { + return ErrInvalid + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return errors.New("begin task resume transaction") + } + defer tx.Rollback() + var accountID, confirmationID string + var confirmationVersion int64 + err = tx.QueryRowContext(ctx, ` + SELECT task.account_id, task.confirmation_id, task.confirmation_version + FROM operation_task task + JOIN social_account account ON account.id = task.account_id + JOIN content_draft draft ON draft.id = task.draft_id + JOIN confirmation confirmation ON confirmation.id = task.confirmation_id + JOIN environment_binding binding ON binding.account_id = task.account_id + JOIN browser_env environment ON environment.alias = binding.browser_env_alias + JOIN network_exit network ON network.id = binding.network_exit_id + JOIN runtime_instance runtime ON runtime.binding_id = binding.id AND runtime.released_at IS NULL + WHERE task.id = $1 AND task.state IN ('policy_hold', 'needs_confirmation') + AND task.verification_result = 'not_executed' + AND account.status = 'active' AND account.authorization_status = 'authorized' + AND account.version = task.account_version + AND draft.account_id = task.account_id AND draft.version = task.draft_version + AND confirmation.account_id = task.account_id AND confirmation.account_version = task.account_version + AND confirmation.draft_id = task.draft_id AND confirmation.draft_version = task.draft_version + AND confirmation.version = task.confirmation_version + AND network.health_status = 'healthy' AND NOT binding.runtime_cleanup_pending + AND runtime.binding_version = binding.version AND runtime.lease_until > now() + FOR UPDATE OF task, account, draft, confirmation, binding, environment, network, runtime`, taskID). + Scan(&accountID, &confirmationID, &confirmationVersion) + if errors.Is(err, sql.ErrNoRows) { + task, taskErr := scanTask(tx.QueryRowContext(ctx, ` + SELECT id, idempotency_key, account_id, account_version, draft_id, draft_version, + confirmation_id, confirmation_version, state, hold_reason, verification_result, verified_at, verified_by, + created_at, updated_at + FROM operation_task WHERE id = $1`, taskID)) + if taskErr != nil { + return taskErr + } + if task.VerificationResult != "not_executed" || (task.State != "policy_hold" && task.State != "needs_confirmation") { + return ErrConflict + } + reason, unavailable, reasonErr := taskReadinessReason(ctx, tx, taskID) + if reasonErr != nil { + return reasonErr + } + if reason == "" { + return ErrConflict + } + return &ReadinessError{Reason: reason, Unavailable: unavailable} + } + if err != nil { + return errors.New("validate task resume readiness") + } + if _, err := tx.ExecContext(ctx, ` + UPDATE operation_task SET state = 'queued', hold_reason = NULL, + verification_result = NULL, verified_at = NULL, verified_by = NULL, + lease_owner = NULL, lease_until = NULL, updated_at = now() + WHERE id = $1`, taskID); err != nil { + return errors.New("resume task") + } + if err := appendAudit(ctx, tx, "task_resumed", "manual_verification_not_executed", accountID, confirmationID, + confirmationVersion, "", taskID, nil); err != nil { + return err + } + return commit(tx) +} + +type rowQuerier interface { + QueryRowContext(context.Context, string, ...any) *sql.Row +} + +func taskReadinessReason(ctx context.Context, queryer rowQuerier, taskID string) (string, bool, error) { + var reason string + err := queryer.QueryRowContext(ctx, ` + SELECT CASE + WHEN account.id IS NULL THEN 'account_missing' + WHEN account.authorization_status = 'revoked' THEN 'account_revoked' + WHEN account.status <> 'active' THEN 'account_paused' + WHEN account.version <> task.account_version THEN 'account_version_changed' + WHEN draft.id IS NULL OR draft.account_id <> task.account_id OR draft.version <> task.draft_version THEN 'draft_version_changed' + WHEN confirmation.id IS NULL THEN 'confirmation_missing' + WHEN confirmation.account_id <> task.account_id OR confirmation.account_version <> task.account_version + OR confirmation.draft_id <> task.draft_id OR confirmation.draft_version <> task.draft_version + OR confirmation.version <> task.confirmation_version THEN 'confirmation_version_changed' + WHEN binding.id IS NULL THEN 'binding_missing' + WHEN environment.alias IS NULL THEN 'environment_missing' + WHEN network.id IS NULL THEN 'exit_missing' + WHEN network.health_status <> 'healthy' THEN 'exit_unhealthy' + WHEN binding.runtime_cleanup_pending THEN 'runtime_stop_pending' + WHEN runtime.id IS NULL THEN 'runtime_missing' + WHEN runtime.binding_version IS DISTINCT FROM binding.version THEN 'binding_version_changed' + WHEN runtime.lease_until <= now() THEN 'runtime_lease_expired' + ELSE '' + END + FROM operation_task task + LEFT JOIN social_account account ON account.id = task.account_id + LEFT JOIN content_draft draft ON draft.id = task.draft_id + LEFT JOIN confirmation confirmation ON confirmation.id = task.confirmation_id + LEFT JOIN environment_binding binding ON binding.account_id = task.account_id + LEFT JOIN browser_env environment ON environment.alias = binding.browser_env_alias + LEFT JOIN network_exit network ON network.id = binding.network_exit_id + LEFT JOIN runtime_instance runtime ON runtime.binding_id = binding.id AND runtime.released_at IS NULL + WHERE task.id = $1`, taskID).Scan(&reason) + if err != nil { + return "", false, rowError(err) + } + unavailable := reason == "binding_missing" || reason == "environment_missing" || reason == "exit_missing" || + reason == "exit_unhealthy" || reason == "runtime_stop_pending" || reason == "runtime_missing" || reason == "runtime_lease_expired" + return reason, unavailable, nil +} + +func (s *Store) FinishTask(ctx context.Context, taskID string) error { + if !refPattern.MatchString(taskID) { + return ErrInvalid + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return errors.New("begin task finish transaction") + } + defer tx.Rollback() + var accountID, result string + var confirmationID sql.NullString + var confirmationVersion sql.NullInt64 + err = tx.QueryRowContext(ctx, ` + UPDATE operation_task SET state = verification_result, hold_reason = NULL, + lease_owner = NULL, lease_until = NULL, updated_at = now() + WHERE id = $1 AND state IN ('policy_hold', 'needs_confirmation') + AND verification_result IN ('succeeded', 'failed') + RETURNING account_id, confirmation_id, confirmation_version, verification_result`, taskID). + Scan(&accountID, &confirmationID, &confirmationVersion, &result) + if err != nil { + return rowError(err) + } + if err := appendAudit(ctx, tx, "task_manually_finished", "manual_verification_"+result, accountID, confirmationID.String, + confirmationVersion.Int64, "", taskID, map[string]string{"state": result}); err != nil { + return err + } + return commit(tx) +} + func (s *Store) CancelTask(ctx context.Context, taskID string) error { if !refPattern.MatchString(taskID) { return ErrInvalid @@ -913,6 +1306,8 @@ func (s *Store) CancelTask(ctx context.Context, taskID string) error { var confirmationVersion sql.NullInt64 err = tx.QueryRowContext(ctx, ` UPDATE operation_task SET state = CASE WHEN state = 'executing' THEN 'needs_confirmation' ELSE 'cancelled' END, + hold_reason = CASE WHEN state = 'executing' THEN 'task_result_uncertain' END, + verification_result = NULL, verified_at = NULL, verified_by = NULL, lease_owner = NULL, lease_until = NULL, updated_at = now() WHERE id = $1 AND state IN ('queued', 'executing', 'needs_confirmation', 'policy_hold') RETURNING account_id, state, current_attempt_id, confirmation_id, confirmation_version`, taskID).Scan( @@ -979,8 +1374,9 @@ func (s *Store) claim(ctx context.Context, workerID string) (Execution, error) { ORDER BY t.created_at, t.id FOR UPDATE OF t, a, binding, network, runtime SKIP LOCKED LIMIT 1 ) - UPDATE operation_task t SET state = 'executing', lease_owner = $1, - lease_until = now() + interval '1 minute', updated_at = now() + UPDATE operation_task t SET state = 'executing', hold_reason = NULL, + verification_result = NULL, verified_at = NULL, verified_by = NULL, + lease_owner = $1, lease_until = now() + interval '1 minute', updated_at = now() FROM candidate WHERE t.id = candidate.id RETURNING t.id, t.account_id, t.confirmation_id, t.confirmation_version`, workerID).Scan( &execution.TaskID, &execution.AccountID, &execution.ConfirmationID, &execution.ConfirmationVersion) @@ -1018,29 +1414,87 @@ func (s *Store) complete(ctx context.Context, workerID string, execution Executi } defer tx.Rollback() var leaseUntil time.Time + var holdReason string err = tx.QueryRowContext(ctx, ` - SELECT lease_until FROM operation_task - WHERE id = $1 AND state = 'executing' AND lease_owner = $2 AND current_attempt_id = $3 - FOR UPDATE`, execution.TaskID, workerID, execution.AttemptID).Scan(&leaseUntil) + SELECT task.lease_until, CASE + WHEN account.id IS NULL THEN 'account_missing' + WHEN account.authorization_status <> 'authorized' THEN 'account_revoked' + WHEN account.status <> 'active' THEN 'account_paused' + WHEN account.version <> task.account_version THEN 'account_version_changed' + WHEN draft.id IS NULL OR draft.account_id <> task.account_id OR draft.version <> task.draft_version THEN 'draft_version_changed' + WHEN confirmation.id IS NULL OR confirmation.account_id <> task.account_id + OR confirmation.account_version <> task.account_version OR confirmation.draft_id <> task.draft_id + OR confirmation.draft_version <> task.draft_version OR confirmation.version <> task.confirmation_version THEN 'confirmation_version_changed' + WHEN binding.id IS NULL THEN 'binding_missing' + WHEN environment.alias IS NULL THEN 'environment_missing' + WHEN network.id IS NULL THEN 'exit_missing' + WHEN network.health_status <> 'healthy' THEN 'exit_unhealthy' + WHEN binding.runtime_cleanup_pending THEN 'runtime_stop_pending' + WHEN runtime.id IS NULL OR runtime.lease_until <= now() THEN 'runtime_missing' + WHEN claim.binding_version IS DISTINCT FROM binding.version + OR claim.browser_env_alias IS DISTINCT FROM environment.alias + OR claim.network_exit_id IS DISTINCT FROM network.id + OR claim.runtime_instance_id IS DISTINCT FROM runtime.id + OR runtime.binding_version IS DISTINCT FROM binding.version THEN 'task_result_uncertain' + ELSE '' + END + FROM operation_task task + LEFT JOIN LATERAL ( + SELECT * FROM social_account WHERE id = task.account_id FOR UPDATE + ) account ON true + LEFT JOIN LATERAL ( + SELECT * FROM content_draft WHERE id = task.draft_id FOR UPDATE + ) draft ON true + LEFT JOIN LATERAL ( + SELECT * FROM confirmation WHERE id = task.confirmation_id FOR UPDATE + ) confirmation ON true + LEFT JOIN LATERAL ( + SELECT * FROM environment_binding WHERE account_id = task.account_id FOR UPDATE + ) binding ON true + LEFT JOIN LATERAL ( + SELECT * FROM browser_env WHERE alias = binding.browser_env_alias FOR UPDATE + ) environment ON true + LEFT JOIN LATERAL ( + SELECT * FROM network_exit WHERE id = binding.network_exit_id FOR UPDATE + ) network ON true + LEFT JOIN LATERAL ( + SELECT * FROM runtime_instance WHERE binding_id = binding.id AND released_at IS NULL FOR UPDATE + ) runtime ON true + LEFT JOIN LATERAL ( + SELECT browser_env_alias, network_exit_id, runtime_instance_id, binding_version + FROM audit_event WHERE attempt_id = $3 AND event_type = 'task_claimed' ORDER BY id DESC LIMIT 1 + ) claim ON true + WHERE task.id = $1 AND task.state = 'executing' AND task.lease_owner = $2 AND task.current_attempt_id = $3 + FOR UPDATE OF task`, execution.TaskID, workerID, execution.AttemptID).Scan(&leaseUntil, &holdReason) if err != nil { return Execution{}, rowError(err) } state := map[string]string{"succeeded": "succeeded", "failed": "failed", "uncertain": "needs_confirmation", "policy_hold": "policy_hold"}[outcome] - if leaseUntil.Before(time.Now()) { + if leaseUntil.Before(time.Now()) || holdReason != "" { outcome, state = "uncertain", "needs_confirmation" } result, _ := json.Marshal(map[string]string{"mock_outcome": outcome}) if _, err := tx.ExecContext(ctx, `UPDATE execution_attempt SET finished_at = now(), outcome = $1, result = $2 WHERE id = $3`, outcome, result, execution.AttemptID); err != nil { return Execution{}, errors.New("finish execution attempt") } + if state == "needs_confirmation" && holdReason == "" { + holdReason = "task_result_uncertain" + } else if state == "policy_hold" { + holdReason = "task_policy_hold" + } if _, err := tx.ExecContext(ctx, ` - UPDATE operation_task SET state = $1, lease_owner = NULL, lease_until = NULL, updated_at = now() WHERE id = $2`, state, execution.TaskID); err != nil { + UPDATE operation_task SET state = $1, hold_reason = NULLIF($2, ''), + verification_result = NULL, verified_at = NULL, verified_by = NULL, + lease_owner = NULL, lease_until = NULL, updated_at = now() WHERE id = $3`, state, holdReason, execution.TaskID); err != nil { return Execution{}, errors.New("finish task") } reason := map[string]string{ "succeeded": "task_succeeded", "failed": "task_failed", "needs_confirmation": "task_result_uncertain", "policy_hold": "task_policy_hold", }[state] + if holdReason != "" && state == "needs_confirmation" { + reason = holdReason + } if err := appendAudit(ctx, tx, "task_finished", reason, execution.AccountID, execution.ConfirmationID, execution.ConfirmationVersion, execution.AttemptID, execution.TaskID, map[string]string{"state": state}); err != nil { return Execution{}, err } @@ -1053,7 +1507,9 @@ func (s *Store) complete(ctx context.Context, workerID string, execution Executi func quarantineExpired(ctx context.Context, tx *sql.Tx) error { rows, err := tx.QueryContext(ctx, ` - UPDATE operation_task SET state = 'needs_confirmation', lease_owner = NULL, lease_until = NULL, updated_at = now() + UPDATE operation_task SET state = 'needs_confirmation', hold_reason = 'execution_lease_expired', + verification_result = NULL, verified_at = NULL, verified_by = NULL, + lease_owner = NULL, lease_until = NULL, updated_at = now() WHERE state = 'executing' AND lease_until < now() RETURNING id, account_id, current_attempt_id, confirmation_id, confirmation_version`) if err != nil { @@ -1145,7 +1601,8 @@ func quarantineInvalid(ctx context.Context, tx *sql.Tx) error { ) FOR UPDATE OF t SKIP LOCKED ) - UPDATE operation_task task SET state = invalid.state, updated_at = now() + UPDATE operation_task task SET state = invalid.state, hold_reason = invalid.reason_code, + verification_result = NULL, verified_at = NULL, verified_by = NULL, updated_at = now() FROM invalid WHERE task.id = invalid.id RETURNING task.id, task.account_id, task.confirmation_id, task.confirmation_version, task.state, invalid.reason_code`) if err != nil { @@ -1178,17 +1635,53 @@ func quarantineInvalid(ctx context.Context, tx *sql.Tx) error { } func (s *Store) Audit(ctx context.Context) ([]AuditEvent, error) { + page, err := s.ListAudit(ctx, AuditFilter{Page: 1, PageSize: 1000}) + return page.Data, err +} + +func (s *Store) ListAudit(ctx context.Context, filter AuditFilter) (AuditPage, error) { + if (filter.AccountID != "" && !idPattern.MatchString(filter.AccountID)) || + (filter.TaskID != "" && !refPattern.MatchString(filter.TaskID)) || + (filter.AttemptID != "" && !refPattern.MatchString(filter.AttemptID)) || + (filter.BrowserEnvAlias != "" && !refPattern.MatchString(filter.BrowserEnvAlias)) || + (filter.NetworkExitID != "" && !refPattern.MatchString(filter.NetworkExitID)) || + (filter.EventType != "" && !eventPattern.MatchString(filter.EventType)) || filter.Page < 1 || + filter.PageSize < 1 || filter.PageSize > 1000 || (filter.From != nil && filter.To != nil && filter.From.After(*filter.To)) { + return AuditPage{}, ErrInvalid + } + var from, to any + if filter.From != nil { + from = *filter.From + } + if filter.To != nil { + to = *filter.To + } + var total int + err := s.db.QueryRowContext(ctx, ` + SELECT count(*) FROM audit_event + WHERE ($1 = '' OR account_id = $1) AND ($2 = '' OR task_id = $2) AND ($3 = '' OR attempt_id = $3) + AND ($4 = '' OR browser_env_alias = $4) AND ($5 = '' OR network_exit_id = $5) AND ($6 = '' OR event_type = $6) + AND ($7::timestamptz IS NULL OR created_at >= $7) AND ($8::timestamptz IS NULL OR created_at <= $8)`, + filter.AccountID, filter.TaskID, filter.AttemptID, filter.BrowserEnvAlias, filter.NetworkExitID, filter.EventType, from, to).Scan(&total) + if err != nil { + return AuditPage{}, errors.New("count audit events") + } rows, err := s.db.QueryContext(ctx, ` SELECT id, event_type, account_id, confirmation_id, confirmation_version, attempt_id, task_id, browser_env_alias, network_exit_id, runtime_instance_id, binding_version, actor, reason_code, operation_id, action, outcome, old_image_version, new_image_version, details, created_at - FROM audit_event ORDER BY id`) + FROM audit_event + WHERE ($1 = '' OR account_id = $1) AND ($2 = '' OR task_id = $2) AND ($3 = '' OR attempt_id = $3) + AND ($4 = '' OR browser_env_alias = $4) AND ($5 = '' OR network_exit_id = $5) AND ($6 = '' OR event_type = $6) + AND ($7::timestamptz IS NULL OR created_at >= $7) AND ($8::timestamptz IS NULL OR created_at <= $8) + ORDER BY created_at DESC, id DESC LIMIT $9 OFFSET $10`, filter.AccountID, filter.TaskID, filter.AttemptID, + filter.BrowserEnvAlias, filter.NetworkExitID, filter.EventType, from, to, filter.PageSize, (filter.Page-1)*filter.PageSize) if err != nil { - return nil, errors.New("read audit events") + return AuditPage{}, errors.New("read audit events") } defer rows.Close() - var events []AuditEvent + events := []AuditEvent{} for rows.Next() { var event AuditEvent var accountID, confirmationID, attemptID, taskID, browserEnvAlias, networkExitID sql.NullString @@ -1198,7 +1691,7 @@ func (s *Store) Audit(ctx context.Context) ([]AuditEvent, error) { &browserEnvAlias, &networkExitID, &runtimeInstanceID, &bindingVersion, &actor, &reasonCode, &operationID, &action, &outcome, &oldImage, &newImage, &event.Details, &event.CreatedAt); err != nil { - return nil, errors.New("decode audit event") + return AuditPage{}, errors.New("decode audit event") } event.AccountID, event.ConfirmationID, event.ConfirmationVersion = accountID.String, confirmationID.String, confirmationVersion.Int64 event.AttemptID, event.TaskID = attemptID.String, taskID.String @@ -1207,9 +1700,50 @@ func (s *Store) Audit(ctx context.Context) ([]AuditEvent, error) { event.Actor, event.ReasonCode = actor.String, reasonCode.String event.OperationID, event.Action, event.Outcome = operationID.String, action.String, outcome.String event.OldImageVersion, event.NewImageVersion = oldImage.String, newImage.String + event.Details = safeDetails(event.Details) events = append(events, event) } - return events, rows.Err() + if err := rows.Err(); err != nil { + return AuditPage{}, errors.New("read audit events") + } + return AuditPage{Data: events, Total: total, Page: filter.Page, PageSize: filter.PageSize}, nil +} + +func safeDetails(raw json.RawMessage) json.RawMessage { + var value any + if json.Unmarshal(raw, &value) != nil { + return json.RawMessage(`{}`) + } + value = allowDetails(value) + encoded, err := json.Marshal(value) + if err != nil { + return json.RawMessage(`{}`) + } + return encoded +} + +var allowedDetailKeys = map[string]bool{ + "platform": true, "draft_id": true, "draft_version": true, "account_version": true, + "tasks_held": true, "attempts_interrupted": true, "verification_result": true, + "state": true, "worker_id": true, "mock_outcome": true, +} + +func allowDetails(value any) any { + switch value := value.(type) { + case map[string]any: + for key, child := range value { + if !allowedDetailKeys[strings.ToLower(key)] { + delete(value, key) + continue + } + value[key] = allowDetails(child) + } + case []any: + for index, child := range value { + value[index] = allowDetails(child) + } + } + return value } func appendAudit(ctx context.Context, tx *sql.Tx, eventType, reasonCode, accountID, confirmationID string, confirmationVersion int64, attemptID, taskID string, details any) error { diff --git a/internal/phasea/store_test.go b/internal/phasea/store_test.go index a3ac050..018e56b 100644 --- a/internal/phasea/store_test.go +++ b/internal/phasea/store_test.go @@ -59,6 +59,31 @@ func TestValidationRejectsInvalidInputsBeforePersistence(t *testing.T) { } } +func TestTaskRecoveryActionsAndEvidenceAreFailClosed(t *testing.T) { + for _, test := range []struct { + task Task + want string + }{ + {task: Task{State: "needs_confirmation", HoldReason: "execution_lease_expired"}, want: "verify"}, + {task: Task{State: "policy_hold", HoldReason: "account_version_changed"}, want: "reconfirm"}, + {task: Task{State: "needs_confirmation", HoldReason: "confirmation_missing"}, want: ""}, + {task: Task{State: "policy_hold", HoldReason: "future_reason"}, want: ""}, + {task: Task{State: "needs_confirmation", HoldReason: "task_result_uncertain", VerificationResult: "not_executed"}, want: "resume"}, + {task: Task{State: "needs_confirmation", HoldReason: "task_result_uncertain", VerificationResult: "succeeded"}, want: "finish"}, + } { + if got := taskAllowedAction(test.task); got != test.want { + t.Fatalf("allowed action: got=%q want=%q task=%+v", got, test.want, test.task) + } + } + redacted := safeDetails(json.RawMessage(`{"state":{"state":"safe","api_key":"a","private_key":"b","credential_key":"c","authorization_header":"d","proxy_url":"e"}}`)) + if string(redacted) != `{"state":{"state":"safe"}}` { + t.Fatalf("unexpected redaction: %s", redacted) + } + if evidence := safeEvidence(json.RawMessage(`{"mock_outcome":"uncertain","cookie":"secret"}`)); len(evidence) != 1 || evidence["mock_outcome"] != "uncertain" { + t.Fatalf("unexpected evidence allowlist: %#v", evidence) + } +} + func TestPhaseAOfflineWorkflow(t *testing.T) { databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") if databaseURL == "" { @@ -276,6 +301,13 @@ func TestPhaseAOfflineWorkflow(t *testing.T) { assertTaskGate(t, store, taskID, test.wantState, test.wantReason) }) } + if err := store.VerifyTask(ctx, "task-gate-exit-unhealthy", "not_executed"); err != nil { + t.Fatal(err) + } + blockedDetail, err := store.GetTaskDetail(ctx, "task-gate-exit-unhealthy") + if err != nil || blockedDetail.AllowedAction != "" || blockedDetail.ReadinessReason != "exit_unhealthy" { + t.Fatalf("unhealthy exit exposed resume after verification: detail=%+v err=%v", blockedDetail, err) + } // A held task is never revived in place: only a newly confirmed task with a new idempotency key may run. if _, err := store.db.ExecContext(ctx, ` @@ -314,6 +346,10 @@ func TestPhaseAOfflineWorkflow(t *testing.T) { } assertCount(t, store, `SELECT count(*) FROM operation_task WHERE id IN ('task-30', 'task-31') AND state = 'needs_confirmation'`, 2) assertCount(t, store, `SELECT count(*) FROM execution_attempt WHERE task_id IN ('task-30', 'task-31')`, 0) + nullConfirmation, err := store.GetTaskDetail(ctx, "task-30") + if err != nil || nullConfirmation.Confirmation.ID != "" || nullConfirmation.AllowedAction != "" { + t.Fatalf("NULL confirmation detail was not readable and fail-closed: detail=%+v err=%v", nullConfirmation, err) + } uncertain := approvedTask(32, "account-b", accountB.Version, "draft-b", "confirmation-b") if _, _, err := store.Enqueue(ctx, uncertain); err != nil { @@ -401,6 +437,146 @@ func TestPhaseAOfflineWorkflow(t *testing.T) { } assertCount(t, store, `SELECT count(*) FROM operation_task WHERE id = 'task-36' AND state = 'needs_confirmation'`, 1) assertCount(t, store, `SELECT count(*) FROM execution_attempt WHERE task_id = 'task-36' AND outcome = 'uncertain'`, 1) + if err := store.ResumeTask(ctx, "task-36"); !errors.Is(err, ErrConflict) { + t.Fatalf("unknown result resumed without verification: %v", err) + } + if err := store.VerifyTask(ctx, "task-36", "not_executed"); err != nil { + t.Fatal(err) + } + if err := store.VerifyTask(ctx, "task-36", "succeeded"); !errors.Is(err, ErrConflict) { + t.Fatalf("repeated verification changed the recorded conclusion: %v", err) + } + verifiedDetail, err := store.GetTaskDetail(ctx, "task-36") + if err != nil || verifiedDetail.AllowedAction != "resume" || verifiedDetail.ReadinessReason != "" { + t.Fatalf("ready verified task did not expose one resume action: detail=%+v err=%v", verifiedDetail, err) + } + if err := store.ResumeTask(ctx, "task-36"); err != nil { + t.Fatal(err) + } + if execution, err := store.ExecuteMock(ctx, "worker-after-verification", "succeeded"); err != nil || !execution.WasClaimed || execution.TaskID != "task-36" { + t.Fatalf("verified task did not resume: execution=%+v err=%v", execution, err) + } + assertCount(t, store, `SELECT count(*) FROM execution_attempt WHERE task_id = 'task-36'`, 2) + attempts, err := store.GetTaskDetail(ctx, "task-36") + if err != nil || len(attempts.Attempts) != 2 { + t.Fatalf("task attempts are not traceable: detail=%+v err=%v", attempts, err) + } + attemptDetail, err := store.GetTaskAttemptDetail(ctx, attempts.Attempts[1].ID) + if err != nil || attemptDetail.TaskID != "task-36" || attemptDetail.BrowserEnvAlias == "" || attemptDetail.NetworkExitID == "" { + t.Fatalf("attempt deep link is incomplete: detail=%+v err=%v", attemptDetail, err) + } + + stale := approvedTask(37, "account-a", accountA.Version, "draft-a", "confirmation-a") + if _, _, err := store.Enqueue(ctx, stale); err != nil { + t.Fatal(err) + } + staleExecution, err := store.claim(ctx, "worker-stale") + if err != nil || staleExecution.TaskID != stale.ID { + t.Fatalf("claim stale completion fixture: execution=%+v err=%v", staleExecution, err) + } + if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'unhealthy' WHERE id = 'exit-shared'`); err != nil { + t.Fatal(err) + } + completed, err := store.complete(ctx, "worker-stale", staleExecution, "succeeded") + if err != nil || completed.State != "needs_confirmation" { + t.Fatalf("stale worker persisted a successful result: execution=%+v err=%v", completed, err) + } + assertCount(t, store, `SELECT count(*) FROM operation_task WHERE id = 'task-37' AND state = 'needs_confirmation' AND hold_reason = 'exit_unhealthy'`, 1) + if _, err := store.db.ExecContext(ctx, `UPDATE network_exit SET health_status = 'healthy' WHERE id = 'exit-shared'`); err != nil { + t.Fatal(err) + } + + concurrent := approvedTask(38, "account-a", accountA.Version, "draft-a", "confirmation-a") + if _, _, err := store.Enqueue(ctx, concurrent); err != nil { + t.Fatal(err) + } + concurrentExecution, err := store.claim(ctx, "worker-concurrent-release") + if err != nil || concurrentExecution.TaskID != concurrent.ID { + t.Fatalf("claim concurrent completion fixture: execution=%+v err=%v", concurrentExecution, err) + } + releaseTx, err := store.db.BeginTx(ctx, nil) + if err != nil { + t.Fatal(err) + } + if _, err := releaseTx.ExecContext(ctx, `UPDATE runtime_instance SET released_at = now() WHERE id = 'runtime-instance-a'`); err != nil { + t.Fatal(err) + } + type completionResult struct { + execution Execution + err error + } + completionStarted := make(chan struct{}) + completionDone := make(chan completionResult, 1) + go func() { + close(completionStarted) + execution, err := store.complete(ctx, "worker-concurrent-release", concurrentExecution, "succeeded") + completionDone <- completionResult{execution: execution, err: err} + }() + <-completionStarted + select { + case result := <-completionDone: + releaseTx.Rollback() + t.Fatalf("completion bypassed an in-flight runtime release: execution=%+v err=%v", result.execution, result.err) + case <-time.After(100 * time.Millisecond): + } + if err := releaseTx.Commit(); err != nil { + t.Fatal(err) + } + select { + case result := <-completionDone: + if result.err != nil || result.execution.State != "needs_confirmation" { + t.Fatalf("completion after runtime release was not quarantined: execution=%+v err=%v", result.execution, result.err) + } + case <-time.After(5 * time.Second): + t.Fatal("completion remained blocked after runtime release committed") + } + assertCount(t, store, `SELECT count(*) FROM operation_task + WHERE id = 'task-38' AND state = 'needs_confirmation' AND hold_reason = 'runtime_missing'`, 1) + concurrentDetail, err := store.GetTaskDetail(ctx, concurrent.ID) + if err != nil || concurrentDetail.RuntimeInstanceID != "runtime-instance-a" { + t.Fatalf("task detail lost its immutable claim runtime: detail=%+v err=%v", concurrentDetail, err) + } + + if _, err := store.db.ExecContext(ctx, ` + INSERT INTO runtime_instance (id, account_id, binding_id, binding_version, runtime_id, lease_until) + VALUES ('runtime-instance-a2', 'account-a', 'binding-a', 1, 'runtime-a2', now() + interval '1 minute')`); err != nil { + t.Fatal(err) + } + replaced := approvedTask(39, "account-a", accountA.Version, "draft-a", "confirmation-a") + if _, _, err := store.Enqueue(ctx, replaced); err != nil { + t.Fatal(err) + } + replacedExecution, err := store.claim(ctx, "worker-replaced-runtime") + if err != nil || replacedExecution.TaskID != replaced.ID { + t.Fatalf("claim replaced runtime fixture: execution=%+v err=%v", replacedExecution, err) + } + if _, err := store.db.ExecContext(ctx, ` + UPDATE runtime_instance SET released_at = now() WHERE id = 'runtime-instance-a2'; + INSERT INTO runtime_instance (id, account_id, binding_id, binding_version, runtime_id, lease_until) + VALUES ('runtime-instance-a3', 'account-a', 'binding-a', 1, 'runtime-a3', now() + interval '1 minute')`); err != nil { + t.Fatal(err) + } + replacedCompletion, err := store.complete(ctx, "worker-replaced-runtime", replacedExecution, "succeeded") + if err != nil || replacedCompletion.State != "needs_confirmation" { + t.Fatalf("old worker completed against a replacement runtime: execution=%+v err=%v", replacedCompletion, err) + } + assertCount(t, store, `SELECT count(*) FROM operation_task + WHERE id = 'task-39' AND state = 'needs_confirmation' AND hold_reason = 'task_result_uncertain'`, 1) + replacedDetail, err := store.GetTaskDetail(ctx, replaced.ID) + if err != nil || replacedDetail.AllowedAction != "verify" || replacedDetail.BrowserEnvAlias != "account-a" || + replacedDetail.NetworkExitID != "exit-shared" || replacedDetail.RuntimeInstanceID != "runtime-instance-a2" || replacedDetail.BindingVersion != 1 { + t.Fatalf("replacement task detail did not preserve the claim snapshot: detail=%+v err=%v", replacedDetail, err) + } + replacedAttempt, err := store.GetTaskAttemptDetail(ctx, replacedExecution.AttemptID) + if err != nil || replacedAttempt.BrowserEnvAlias != "account-a" || replacedAttempt.NetworkExitID != "exit-shared" || + replacedAttempt.RuntimeInstanceID != "runtime-instance-a2" || replacedAttempt.BindingVersion != 1 { + t.Fatalf("replacement attempt detail did not preserve the claim snapshot: detail=%+v err=%v", replacedAttempt, err) + } + + if _, err := store.db.ExecContext(ctx, `INSERT INTO audit_event (event_type, details) VALUES ('redaction_test', + '{"state":"safe","api_key":"api-secret","private_key":"private-secret","credential_key":"credential-secret","authorization_header":"auth-secret","proxy_url":"proxy-secret","nested":{"state":"nested-safe","token":"nested-secret"}}')`); err != nil { + t.Fatal(err) + } events, err := store.Audit(ctx) if err != nil { @@ -417,11 +593,19 @@ func TestPhaseAOfflineWorkflow(t *testing.T) { t.Fatal("audit does not trace confirmation version, task, and attempt") } exported, _ := json.Marshal(events) - for _, forbidden := range []string{"password", "cookie", "token", "credential-a", "creatorhub/account-a"} { + for _, forbidden := range []string{"password", "cookie", "token", "credential-a", "creatorhub/account-a", "api-secret", "private-secret", "credential-secret", "auth-secret", "proxy-secret", "nested-secret"} { if strings.Contains(strings.ToLower(string(exported)), forbidden) { t.Fatalf("audit export contains sensitive field or credential reference %q", forbidden) } } + firstPage, err := store.ListAudit(ctx, AuditFilter{Page: 1, PageSize: 1}) + if err != nil || firstPage.Total < 2 || len(firstPage.Data) != 1 { + t.Fatalf("first audit page: page=%+v err=%v", firstPage, err) + } + secondPage, err := store.ListAudit(ctx, AuditFilter{Page: 2, PageSize: 1}) + if err != nil || len(secondPage.Data) != 1 || secondPage.Data[0].ID == firstPage.Data[0].ID { + t.Fatalf("second audit page: page=%+v err=%v", secondPage, err) + } if _, err := store.db.ExecContext(ctx, `UPDATE audit_event SET event_type = 'rewritten' WHERE id = 1`); err == nil { t.Fatal("audit events must be append-only") } @@ -548,7 +732,10 @@ func applyHubMigrationsForPhaseATest(t *testing.T, store *Store) { for _, migrationFile := range []struct { version int name string - }{{2, "002_hub.sql"}, {3, "003_unified_accounts.sql"}, {4, "004_environment_actions.sql"}, {5, "005_sanitize_legacy_proxy.sql"}, {6, "006_runtime_cleanup.sql"}, {7, "007_runtime_binding_version.sql"}} { + }{{2, "002_hub.sql"}, {3, "003_unified_accounts.sql"}, {4, "004_environment_actions.sql"}, {5, "005_sanitize_legacy_proxy.sql"}, + {6, "006_runtime_cleanup.sql"}, {7, "007_runtime_binding_version.sql"}, {8, "008_runtime_cleanup_generation.sql"}, + {9, "009_runtime_cleanup_compatibility.sql"}, {10, "010_runtime_network_generation.sql"}, {11, "011_task_recovery.sql"}, + {12, "012_task_recovery_compatibility.sql"}} { var applied bool if err := store.db.QueryRow(`SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migrationFile.version).Scan(&applied); err != nil { t.Fatal(err) diff --git a/web/src/AuditList.jsx b/web/src/AuditList.jsx new file mode 100644 index 0000000..332964a --- /dev/null +++ b/web/src/AuditList.jsx @@ -0,0 +1,91 @@ +import { useEffect, useState } from 'react' +import { useGetList } from 'ra-core' +import { Link as RouterLink, useSearchParams } from 'react-router-dom' +import { + Alert, + Box, + Button, + CircularProgress, + Paper, + Stack, + Table, + TableBody, + TableCell, + TableContainer, + TableHead, + TableRow, + TextField, + Typography, +} from '@mui/material' +import HistoryOutlined from '@mui/icons-material/HistoryOutlined' + +const wrapAnywhere = { overflowWrap: 'anywhere', minWidth: 0 } +const dayStart = value => value ? new Date(`${value}T00:00:00`).toISOString() : '' +const dayEnd = value => value ? new Date(`${value}T23:59:59.999`).toISOString() : '' + +export function auditCSV(events) { + const fields = ['id', 'event_type', 'account_id', 'task_id', 'attempt_id', 'browser_env_alias', 'network_exit_id', 'binding_version', 'reason_code', 'outcome', 'created_at'] + const quote = value => `"${String(value ?? '').replaceAll('"', '""')}"` + return [fields, ...events.map(event => fields.map(field => event[field]))].map(row => row.map(quote).join(',')).join('\n') +} + +function exportPage(events) { + const url = URL.createObjectURL(new Blob([auditCSV(events)], { type: 'text/csv;charset=utf-8' })) + const link = document.createElement('a') + link.href = url + link.download = 'creatorhub-audit.csv' + link.click() + URL.revokeObjectURL(url) +} + +function AuditLinks({ event }) { + return ( + + {event.account_id ? 账号:{event.account_id} : null} + {event.task_id ? 任务:{event.task_id} : null} + {event.attempt_id ? attempt:{event.attempt_id} : null} + + ) +} + +export function AuditList() { + const [searchParams] = useSearchParams() + const [accountID, setAccountID] = useState(searchParams.get('account_id') || '') + const [taskID, setTaskID] = useState(searchParams.get('task_id') || '') + const [attemptID, setAttemptID] = useState(searchParams.get('attempt_id') || '') + const [browserEnvAlias, setBrowserEnvAlias] = useState(searchParams.get('browser_env_alias') || '') + const [networkExitID, setNetworkExitID] = useState(searchParams.get('network_exit_id') || '') + const [eventType, setEventType] = useState('') + const [from, setFrom] = useState('') + const [to, setTo] = useState('') + const [page, setPage] = useState(1) + const perPage = 25 + const filter = { account_id: accountID, task_id: taskID, attempt_id: attemptID, browser_env_alias: browserEnvAlias, network_exit_id: networkExitID, event_type: eventType, from: dayStart(from), to: dayEnd(to) } + const { data: events = [], total = 0, error, isPending } = useGetList('audit', { pagination: { page, perPage }, filter }, { retry: false }) + + useEffect(() => { document.title = 'CreatorHub · 审计记录' }, []) + useEffect(() => { setPage(1) }, [accountID, taskID, attemptID, browserEnvAlias, networkExitID, eventType, from, to]) + + return ( + + 审计记录追加式事件回溯;页面不提供编辑或删除 + + setAccountID(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '按账号筛选审计' } }} /> + setTaskID(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '按任务筛选审计' } }} /> + setAttemptID(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '按 attempt 筛选审计' } }} /> + setBrowserEnvAlias(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '按环境筛选审计' } }} /> + setNetworkExitID(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '按出口筛选审计' } }} /> + setEventType(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '按事件筛选审计' } }} /> + setFrom(event.target.value)} slotProps={{ inputLabel: { shrink: true }, htmlInput: { 'aria-label': '审计开始日期' } }} /> + setTo(event.target.value)} slotProps={{ inputLabel: { shrink: true }, htmlInput: { 'aria-label': '审计结束日期' } }} /> + + {accountID || taskID || attemptID || browserEnvAlias || networkExitID || eventType || from || to ? : null} + {error ? {error.message} : null} + {isPending ? : null} + {!isPending && events.length === 0 ? {accountID || taskID || attemptID || browserEnvAlias || networkExitID || eventType || from || to ? '没有符合筛选条件的审计事件' : '尚无审计事件'} : null} + {events.length ? 事件关联环境 / 出口原因 / 时间{events.map(event => {event.event_type}#{event.id}{event.browser_env_alias ? {event.browser_env_alias} : '—'} / {event.network_exit_id ? {event.network_exit_id} : '—'}绑定版本 {event.binding_version || '—'}{event.reason_code || event.outcome || '—'}{new Date(event.created_at).toLocaleString('zh-CN')})}
: null} + {events.length ? {events.map(event => {event.event_type} · #{event.id}环境 / 出口:{event.browser_env_alias ? {event.browser_env_alias} : '—'} / {event.network_exit_id ? {event.network_exit_id} : '—'} · 绑定版本 {event.binding_version || '—'}原因:{event.reason_code || event.outcome || '—'}{new Date(event.created_at).toLocaleString('zh-CN')})} : null} + {total > perPage ? 第 {page} / {Math.ceil(total / perPage)} 页 : null} +
+ ) +} diff --git a/web/src/AuditList.test.jsx b/web/src/AuditList.test.jsx new file mode 100644 index 0000000..45eef0e --- /dev/null +++ b/web/src/AuditList.test.jsx @@ -0,0 +1,51 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { cleanup, fireEvent, render, screen, waitFor } from '@testing-library/react' +import { CoreAdminContext } from 'ra-core' +import { MemoryRouter } from 'react-router-dom' +import { AuditList, auditCSV } from './AuditList' + +afterEach(() => { cleanup(); vi.restoreAllMocks() }) + +describe('AuditList', () => { + it('deep-links safe correlation fields and exposes no mutation or secret details', async () => { + const dataProvider = { + getList: vi.fn().mockResolvedValue({ data: [{ + id: 7, event_type: 'task_verified', account_id: 'account-a', task_id: 'task-a', attempt_id: 'attempt-a', + browser_env_alias: 'env-a', network_exit_id: 'exit-a', binding_version: 4, reason_code: 'manual_verification_recorded', + details: { token: 'must-not-render' }, created_at: '2026-08-30T00:00:00Z', + }], total: 1 }), + getOne: vi.fn(), getMany: vi.fn(), getManyReference: vi.fn(), create: vi.fn(), update: vi.fn(), updateMany: vi.fn(), delete: vi.fn(), deleteMany: vi.fn(), + } + render() + + expect((await screen.findAllByRole('link', { name: 'task-a' })).length).toBeGreaterThan(0) + expect(screen.getAllByRole('link', { name: 'account-a' })[0].getAttribute('href')).toBe('/accounts/account-a') + expect(screen.getAllByRole('link', { name: 'attempt-a' })[0].getAttribute('href')).toBe('/attempts/attempt-a') + expect(screen.getAllByRole('link', { name: 'env-a' })[0].getAttribute('href')).toBe('/browsers/env-a') + expect(screen.getAllByRole('link', { name: 'exit-a' })[0].getAttribute('href')).toBe('/network-exits/exit-a') + expect(screen.queryByText('must-not-render')).toBeNull() + expect(screen.queryByRole('button', { name: /编辑|删除/ })).toBeNull() + }) + + it('exports only allowlisted correlation columns and requests the next page', async () => { + const event = { + id: 7, event_type: 'task_verified', task_id: 'task-a', attempt_id: 'attempt-a', browser_env_alias: 'env-a', network_exit_id: 'exit-a', + details: { api_key: 'api-secret', private_key: 'private-secret', credential_key: 'credential-secret', authorization_header: 'auth-secret', proxy_url: 'proxy-secret' }, + created_at: '2026-08-30T00:00:00Z', + } + const csv = auditCSV([event]) + expect(csv).toContain('attempt-a') + for (const secret of ['api-secret', 'private-secret', 'credential-secret', 'auth-secret', 'proxy-secret']) expect(csv).not.toContain(secret) + + const getList = vi.fn().mockImplementation((_resource, params) => Promise.resolve({ + data: [{ ...event, id: params.pagination.page }], total: 26, + })) + const dataProvider = { + getList, getOne: vi.fn(), getMany: vi.fn(), getManyReference: vi.fn(), create: vi.fn(), update: vi.fn(), updateMany: vi.fn(), delete: vi.fn(), deleteMany: vi.fn(), + } + render() + + fireEvent.click(await screen.findByRole('button', { name: '下一页' })) + await waitFor(() => expect(getList.mock.calls.some(([, params]) => params.pagination.page === 2)).toBe(true)) + }) +}) diff --git a/web/src/TaskCenter.jsx b/web/src/TaskCenter.jsx new file mode 100644 index 0000000..f078569 --- /dev/null +++ b/web/src/TaskCenter.jsx @@ -0,0 +1,148 @@ +import { useEffect, useState } from 'react' +import { useDataProvider, useGetList, useGetOne } from 'ra-core' +import { Link as RouterLink, useParams } from 'react-router-dom' +import { + Alert, + Box, + Button, + CircularProgress, + MenuItem, + Paper, + Stack, + Table, + TableBody, + TableCell, + TableContainer, + TableHead, + TableRow, + TextField, + Typography, +} from '@mui/material' +import AssignmentTurnedInOutlined from '@mui/icons-material/AssignmentTurnedInOutlined' + +const wrapAnywhere = { overflowWrap: 'anywhere', minWidth: 0 } + +export const taskLabels = { + queued: '已排队', + executing: '执行中', + succeeded: '已成功', + failed: '失败', + needs_confirmation: '需要人工确认', + policy_hold: '策略暂停', + cancelled: '已取消', +} + +export const reasonLabels = { + account_paused: '账号已暂停', + account_revoked: '账号授权已撤销', + binding_missing: '运行环境绑定缺失', + environment_missing: '原运行环境不可用', + exit_missing: '固定出口缺失', + exit_unhealthy: '固定出口失败或发生漂移', + runtime_stop_pending: '停止结果未知,等待人工核验', + runtime_missing: '原运行实例不可用', + runtime_lease_expired: '运行实例租约已过期', + execution_lease_expired: '执行租约过期,结果未知', + task_result_uncertain: '执行结果不确定', + task_policy_hold: '执行器请求策略暂停', + account_version_changed: '账号版本已变化', + draft_version_changed: '草稿版本已变化', + confirmation_missing: '确认快照缺失', + confirmation_version_changed: '确认版本已变化', + binding_version_changed: '环境或出口版本已变化', + legacy_confirmation_required: '历史任务结果未知,等待人工核验', +} + +const stateText = state => taskLabels[state] || `未知状态:${state}` +const reasonText = reason => reason ? (reasonLabels[reason] || `未知停止原因:${reason}`) : '无' + +function TaskSummary({ task }) { + return ( + <> + {task.id} + {stateText(task.state)} + {reasonText(task.hold_reason)} + + ) +} + +export function TaskList() { + const [accountID, setAccountID] = useState('') + const [state, setState] = useState('') + const { data: tasks = [], error, isPending } = useGetList('tasks', { filter: { account_id: accountID, state } }, { retry: false }) + + useEffect(() => { document.title = 'CreatorHub · 任务中心' }, []) + + return ( + + + 任务中心 + 逐条处理暂停、版本变化与未知执行结果;不会自动重试或更换出口 + + + setAccountID(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '按账号筛选' } }} /> + setState(event.target.value)} sx={{ minWidth: 220 }} slotProps={{ htmlInput: { 'aria-label': '按状态筛选' } }}> + 全部状态 + {Object.entries(taskLabels).map(([value, label]) => {label})} + + + {error ? {error.message} : null} + {isPending ? : null} + {!isPending && tasks.length === 0 ? {accountID || state ? '没有符合筛选条件的任务' : '尚无任务'}{accountID || state ? : null} : null} + {tasks.length > 0 ? 任务账号确认更新时间{tasks.map(task => {task.account_id}{task.confirmation_id || '未确认'}{new Date(task.updated_at || task.created_at).toLocaleString('zh-CN')})}
: null} + {tasks.length > 0 ? {tasks.map(task => 账号:{task.account_id}确认:{task.confirmation_id || '未确认'}{new Date(task.updated_at || task.created_at).toLocaleString('zh-CN')})} : null} +
+ ) +} + +function actionError(error) { + const reason = error?.body?.reason_code + if (error?.status === 409) return `当前版本仍不允许恢复(409):${reasonText(reason)}` + if (error?.status === 503) return `原资源尚未恢复(503):${reasonText(reason)}` + return error?.message || '操作失败' +} + +export function TaskDetail() { + const { id } = useParams() + const dataProvider = useDataProvider() + const [verification, setVerification] = useState('not_executed') + const [busy, setBusy] = useState(false) + const [message, setMessage] = useState(null) + const { data: task, error, isPending, refetch } = useGetOne('tasks', { id }, { retry: false }) + + useEffect(() => { document.title = 'CreatorHub · 任务详情' }, []) + + if (isPending) return + if (error || !task) return {error?.message || '任务不存在'} + + async function act(action, data) { + setBusy(true) + setMessage(null) + try { + await dataProvider.taskAction(task.id, action, data) + setMessage({ severity: 'success', text: action === 'verify' ? '人工核验已记录;现在只能按该结论恢复或结束。' : action === 'resume' ? '任务已恢复排队。' : '任务已按核验结论结束。' }) + await refetch() + } catch (reason) { + setMessage({ severity: 'error', text: actionError(reason) }) + } finally { + setBusy(false) + } + } + + const confirmation = task.confirmation || {} + const unknownReason = task.hold_reason && !reasonLabels[task.hold_reason] + return ( + + 任务详情{task.id} + {message ? {message.text} : null} + {unknownReason ? 后端返回了未知停止原因“{task.hold_reason}”;为避免不安全操作,本页不提供恢复动作。 : null} + + 当前状态{stateText(task.state)}停止原因:{reasonText(task.hold_reason)}人工核验:{task.verification_result || '尚未记录'} + 确认快照确认:{confirmation.id || '未记录'} · v{confirmation.version || '—'}账号版本 {confirmation.account_version || '—'} · 草稿版本 {confirmation.draft_version || '—'}环境:{confirmation.browser_env_alias || '未记录'} · 出口:{confirmation.network_exit_id || '未记录'} · 绑定版本:{confirmation.binding_version || '未记录'} + 执行关联环境:{task.browser_env_alias || '未记录'}固定出口:{task.network_exit_id || '未记录'}运行实例:{task.runtime_instance_id || '未记录'} · 绑定版本:{task.binding_version || '未记录'}{task.browser_env_alias ? : null}{task.network_exit_id ? : null} + 唯一安全动作{task.readiness_reason ? 当前仍被阻断:{reasonText(task.readiness_reason)}。修复原资源并刷新后,才会开放恢复。 : null}{task.allowed_action === 'verify' ? <>先在原环境和原出口核验真实结果。选择“未执行”不会立即重试,仍需下一步显式恢复。 setVerification(event.target.value)} slotProps={{ htmlInput: { 'aria-label': '人工核验结论' } }}>确认未执行,可评估恢复确认已成功确认已失败 : null}{task.allowed_action === 'resume' ? : null}{task.allowed_action === 'finish' ? : null}{task.allowed_action === 'reconfirm' ? <>版本已变化,旧确认不可继续使用。 : null}{!task.allowed_action ? 当前没有可安全执行的动作。 : null} + + 执行尝试{task.attempts?.length ? {task.attempts.map(attempt => {attempt.id}{attempt.outcome || '执行中'} · {new Date(attempt.started_at).toLocaleString('zh-CN')}脱敏证据:{attempt.evidence?.mock_outcome || '无可公开证据'})} : 尚未产生执行尝试。} + + ) +} diff --git a/web/src/TaskCenter.test.jsx b/web/src/TaskCenter.test.jsx new file mode 100644 index 0000000..5bb04e2 --- /dev/null +++ b/web/src/TaskCenter.test.jsx @@ -0,0 +1,67 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { cleanup, fireEvent, render, screen, waitFor } from '@testing-library/react' +import { CoreAdminContext } from 'ra-core' +import { MemoryRouter, Route, Routes } from 'react-router-dom' +import { TaskDetail } from './TaskCenter' + +afterEach(() => { cleanup(); vi.restoreAllMocks() }) + +const task = { + id: 'task-a', account_id: 'account-a', draft_id: 'draft-a', confirmation_id: 'confirmation-a', state: 'needs_confirmation', + hold_reason: 'task_result_uncertain', allowed_action: 'verify', browser_env_alias: 'env-a', network_exit_id: 'exit-a', binding_version: 4, + confirmation: { id: 'confirmation-a', version: 1, account_version: 2, draft_version: 1, browser_env_alias: 'env-a', network_exit_id: 'exit-a', binding_version: 4 }, + attempts: [{ id: 'attempt-a', started_at: '2026-08-30T00:00:00Z', outcome: 'uncertain', evidence: { mock_outcome: 'uncertain', token: 'must-not-render' } }], +} + +function provider(record = task) { + return { + getOne: vi.fn().mockResolvedValue({ data: record }), + taskAction: vi.fn().mockResolvedValue(), + getList: vi.fn(), getMany: vi.fn(), getManyReference: vi.fn(), create: vi.fn(), update: vi.fn(), updateMany: vi.fn(), delete: vi.fn(), deleteMany: vi.fn(), + } +} + +function renderTask(dataProvider) { + return render(} />) +} + +describe('TaskDetail', () => { + it('requires a structured manual verification before any resume', async () => { + const dataProvider = provider() + renderTask(dataProvider) + + expect(await screen.findByText(/执行结果不确定/)).toBeTruthy() + expect(screen.queryByRole('button', { name: /恢复排队/ })).toBeNull() + fireEvent.click(screen.getByRole('button', { name: '记录人工核验' })) + await waitFor(() => expect(dataProvider.taskAction).toHaveBeenCalledWith('task-a', 'verify', { result: 'not_executed' })) + expect(screen.queryByText('must-not-render')).toBeNull() + expect(screen.getByRole('link', { name: 'attempt-a' }).getAttribute('href')).toBe('/attempts/attempt-a') + expect(screen.getByRole('link', { name: '查看原环境' }).getAttribute('href')).toBe('/browsers/env-a') + expect(screen.getByRole('link', { name: '查看原出口' }).getAttribute('href')).toBe('/network-exits/exit-a') + }) + + it('shows only reconfirm for changed versions', async () => { + renderTask(provider({ ...task, hold_reason: 'account_version_changed', allowed_action: 'reconfirm', attempts: [] })) + + expect((await screen.findByRole('link', { name: '重新核对并确认草稿' })).getAttribute('href')).toBe('/drafts/draft-a') + expect(screen.queryByRole('button', { name: '记录人工核验' })).toBeNull() + expect(screen.queryByRole('button', { name: /恢复排队/ })).toBeNull() + }) + + it('round-trips unknown enums without offering an unsafe action', async () => { + renderTask(provider({ ...task, state: 'future_state', hold_reason: 'future_reason', allowed_action: '', attempts: [] })) + + expect(await screen.findByText('未知状态:future_state')).toBeTruthy() + expect(screen.getByText(/未知停止原因“future_reason”/)).toBeTruthy() + expect(screen.queryByRole('button', { name: '记录人工核验' })).toBeNull() + expect(screen.queryByRole('button', { name: /恢复排队|结束任务/ })).toBeNull() + }) + + it('renders a NULL confirmation without exposing recovery or reconfirm actions', async () => { + renderTask(provider({ ...task, confirmation_id: '', confirmation: {}, hold_reason: 'confirmation_missing', allowed_action: '', attempts: [] })) + + expect(await screen.findByText(/确认:未记录/)).toBeTruthy() + expect(screen.queryByRole('link', { name: '重新核对并确认草稿' })).toBeNull() + expect(screen.queryByRole('button', { name: /恢复排队|记录人工核验|结束任务/ })).toBeNull() + }) +}) diff --git a/web/src/TraceDetails.jsx b/web/src/TraceDetails.jsx new file mode 100644 index 0000000..f236b53 --- /dev/null +++ b/web/src/TraceDetails.jsx @@ -0,0 +1,33 @@ +import { useEffect } from 'react' +import { useGetOne } from 'ra-core' +import { Alert, Box, Button, CircularProgress, Paper, Stack, Typography } from '@mui/material' +import { Link as RouterLink, useParams } from 'react-router-dom' + +const wrapAnywhere = { overflowWrap: 'anywhere', minWidth: 0 } + +function DetailState({ pending, error, children }) { + if (pending) return + if (error) return {error.message} + return children +} + +export function AttemptDetail() { + const { id } = useParams() + const { data: attempt, error, isPending } = useGetOne('attempts', { id }, { retry: false }) + useEffect(() => { document.title = 'CreatorHub · Attempt 详情' }, []) + return {attempt ? Attempt 详情{attempt.id}结果:{attempt.outcome || '执行中'}开始:{new Date(attempt.started_at).toLocaleString('zh-CN')}环境:{attempt.browser_env_alias || '未记录'} · 出口:{attempt.network_exit_id || '未记录'} · 绑定版本:{attempt.binding_version || '未记录'}{attempt.browser_env_alias ? : null}{attempt.network_exit_id ? : null} : null} +} + +export function BrowserDetail() { + const { id } = useParams() + const { data: environment, error, isPending } = useGetOne('browsers', { id }, { retry: false }) + useEffect(() => { document.title = 'CreatorHub · 环境详情' }, []) + return {environment ? 运行环境详情{environment.name}({environment.alias})账号:{environment.account_id} · 镜像:{environment.image_version}固定出口:{environment.network_exit?.id || '未绑定'} · 健康:{environment.network_exit?.health_status || '未知'}绑定版本:{environment.binding_version} · 运行实例:{environment.runtime_instance_id || '无'}{environment.network_exit?.id ? : null} : null} +} + +export function NetworkExitDetail() { + const { id } = useParams() + const { data: exit, error, isPending } = useGetOne('network-exits', { id }, { retry: false }) + useEffect(() => { document.title = 'CreatorHub · 出口详情' }, []) + return {exit ? 网络出口详情{exit.id}{exit.protocol}://{exit.host}:{exit.port}健康:{exit.health_status} · 版本:{exit.version}凭据引用:{exit.credential_reference?.id || '无'} : null} +} diff --git a/web/src/dataProvider.js b/web/src/dataProvider.js index 9656bf6..fb1d6e7 100644 --- a/web/src/dataProvider.js +++ b/web/src/dataProvider.js @@ -25,6 +25,8 @@ const resourcePaths = { drafts: '/phase-a/drafts', confirmations: '/phase-a/confirmations', tasks: '/phase-a/tasks', + attempts: '/phase-a/attempts', + audit: '/phase-a/audit', 'network-exits': '/network-exits', } @@ -32,10 +34,18 @@ export const dataProvider = { async getList(resource, params = {}) { const path = resourcePaths[resource] if (!path) return unsupported(resource, 'getList') - const filterKeys = { drafts: ['account_id'], confirmations: ['draft_id'], tasks: ['account_id', 'draft_id'] }[resource] || [] + const filterKeys = { + drafts: ['account_id'], confirmations: ['draft_id'], tasks: ['account_id', 'draft_id', 'state'], + audit: ['account_id', 'task_id', 'attempt_id', 'browser_env_alias', 'network_exit_id', 'event_type', 'from', 'to'], + }[resource] || [] const query = new URLSearchParams(filterKeys.flatMap(key => params.filter?.[key] ? [[key, params.filter[key]]] : [])) + if (resource === 'audit') { + query.set('page', params.pagination?.page || 1) + query.set('page_size', params.pagination?.perPage || 25) + } const records = await request(`${path}${query.size ? `?${query}` : ''}`) - return { data: records.map(record => ({ ...record, id: record.id ?? record.alias ?? record.version ?? record.name })), total: records.length } + const data = Array.isArray(records) ? records : records.data + return { data: data.map(record => ({ ...record, id: record.id ?? record.alias ?? record.version ?? record.name })), total: records.total ?? data.length } }, async create(resource, { data }) { const path = resourcePaths[resource] @@ -59,7 +69,7 @@ export const dataProvider = { }, async getOne(resource, { id }) { const path = resourcePaths[resource] - if (!path || !['accounts', 'network-exits', 'drafts', 'confirmations', 'tasks'].includes(resource)) return unsupported(resource, 'getOne') + if (!path || !['accounts', 'network-exits', 'browsers', 'drafts', 'confirmations', 'tasks', 'attempts'].includes(resource)) return unsupported(resource, 'getOne') const record = await request(`${path}/${encodeURIComponent(id)}`) return { data: { ...record, id: record.id ?? id } } }, @@ -99,6 +109,11 @@ export const dataProvider = { enqueueConfirmation(confirmationID) { return request('/phase-a/tasks', jsonOptions('POST', { confirmation_id: confirmationID })) }, + async taskAction(id, action, data) { + if (!['verify', 'resume', 'finish', 'cancel'].includes(action)) throw new Error(`未知任务操作: ${action}`) + const options = data ? jsonOptions('POST', data) : { method: 'POST' } + await request(`/phase-a/tasks/${encodeURIComponent(id)}/${action}`, options) + }, async networkExitAction(id, action) { if (action !== 'check' && action !== 'disable') throw new Error(`未知网络出口操作: ${action}`) return request(`/network-exits/${encodeURIComponent(id)}/${action}`, { method: 'POST' }) diff --git a/web/src/dataProvider.test.js b/web/src/dataProvider.test.js index f2b518a..aed9a53 100644 --- a/web/src/dataProvider.test.js +++ b/web/src/dataProvider.test.js @@ -29,9 +29,11 @@ describe('dataProvider', () => { it.each([ ['accounts', 'account-a', '/api/phase-a/accounts/account-a'], ['network-exits', 'exit/one', '/api/network-exits/exit%2Fone'], + ['browsers', 'environment-one', '/api/browsers/environment-one'], ['drafts', 'draft/one', '/api/phase-a/drafts/draft%2Fone'], ['confirmations', 'confirmation/one', '/api/phase-a/confirmations/confirmation%2Fone'], ['tasks', 'task/one', '/api/phase-a/tasks/task%2Fone'], + ['attempts', 'attempt/one', '/api/phase-a/attempts/attempt%2Fone'], ])('loads %s detail through its stable API path', async (resource, id, path) => { const fetch = vi.fn().mockResolvedValue(new Response(JSON.stringify({ id }), { status: 200 })) vi.stubGlobal('fetch', fetch) @@ -57,6 +59,16 @@ describe('dataProvider', () => { expect(fetch).toHaveBeenCalledWith('/api/phase-a/drafts?account_id=account-a', undefined) }) + it('maps paginated audit filters without exposing a CRUD mutation', async () => { + const fetch = vi.fn().mockResolvedValue(new Response('{"data":[{"id":7,"event_type":"task_verified"}],"total":31}', { status: 200 })) + vi.stubGlobal('fetch', fetch) + + await expect(dataProvider.getList('audit', { pagination: { page: 2, perPage: 25 }, filter: { task_id: 'task-a', attempt_id: 'attempt-a', browser_env_alias: 'env-a', network_exit_id: 'exit-a' } })).resolves.toEqual({ + data: [{ id: 7, event_type: 'task_verified' }], total: 31, + }) + expect(fetch).toHaveBeenCalledWith('/api/phase-a/audit?task_id=task-a&attempt_id=attempt-a&browser_env_alias=env-a&network_exit_id=exit-a&page=2&page_size=25', undefined) + }) + it('keeps generated draft, confirmation and enqueue identifiers out of user input', async () => { const fetch = vi.fn().mockImplementation(() => Promise.resolve(new Response('{"id":"generated"}', { status: 201 }))) vi.stubGlobal('fetch', fetch) @@ -108,4 +120,17 @@ describe('dataProvider', () => { await dataProvider.accountAction('account-a', action) expect(fetch).toHaveBeenCalledWith(path, { method: 'POST' }) }) + + it('keeps verification and recovery as separate task actions', async () => { + const fetch = vi.fn().mockResolvedValue(new Response(null, { status: 204 })) + vi.stubGlobal('fetch', fetch) + + await dataProvider.taskAction('task/a', 'verify', { result: 'not_executed' }) + await dataProvider.taskAction('task/a', 'resume') + + expect(fetch.mock.calls).toEqual([ + ['/api/phase-a/tasks/task%2Fa/verify', expect.objectContaining({ method: 'POST', body: '{"result":"not_executed"}' })], + ['/api/phase-a/tasks/task%2Fa/resume', { method: 'POST' }], + ]) + }) }) diff --git a/web/src/layout.jsx b/web/src/layout.jsx index d0738ef..30b0ce0 100644 --- a/web/src/layout.jsx +++ b/web/src/layout.jsx @@ -15,6 +15,7 @@ import { import ChevronLeft from '@mui/icons-material/ChevronLeft' import ChevronRight from '@mui/icons-material/ChevronRight' import AccountCircleOutlined from '@mui/icons-material/AccountCircleOutlined' +import AssignmentTurnedInOutlined from '@mui/icons-material/AssignmentTurnedInOutlined' import DnsOutlined from '@mui/icons-material/DnsOutlined' import HistoryOutlined from '@mui/icons-material/HistoryOutlined' import HubOutlined from '@mui/icons-material/HubOutlined' @@ -40,7 +41,11 @@ export function CreatorHubMenu({ collapsed, onNavigate }) { {collapsed ? null : } - + + + {collapsed ? null : } + + {collapsed ? null : } diff --git a/web/src/main.jsx b/web/src/main.jsx index 6234aac..beea254 100644 --- a/web/src/main.jsx +++ b/web/src/main.jsx @@ -5,6 +5,7 @@ import { CoreAdmin, CustomRoutes, Resource } from 'ra-core' import { CssBaseline, ThemeProvider } from '@mui/material' import { Route } from 'react-router-dom' import { AccountDetail, AccountList } from './AccountList' +import { AuditList } from './AuditList' import { BrowserImageList } from './BrowserImageList' import { BrowserList } from './BrowserList' import { dataProvider } from './dataProvider' @@ -12,6 +13,8 @@ import { DraftDetail } from './DraftDetail' import { GatewayList } from './GatewayList' import { CreatorHubLayout } from './layout' import { NetworkExitList } from './NetworkExitList' +import { TaskDetail, TaskList } from './TaskCenter' +import { AttemptDetail, BrowserDetail, NetworkExitDetail } from './TraceDetails' import { theme } from './theme' import './styles.css' @@ -21,6 +24,8 @@ createRoot(document.getElementById('root')).render( + + @@ -28,6 +33,10 @@ createRoot(document.getElementById('root')).render( } /> } /> + } /> + } /> + } /> + } /> diff --git a/web/tests/responsive.e2e.js b/web/tests/responsive.e2e.js index 225cc7f..d01f0b6 100644 --- a/web/tests/responsive.e2e.js +++ b/web/tests/responsive.e2e.js @@ -171,3 +171,68 @@ test('keeps draft review responsive and restores focus after dialog close and su await expect(page.getByRole('dialog')).toBeHidden() await expect(trigger).toBeFocused() }) + +test('keeps task recovery and audit traceable without blind retry at 599px, 900px and 1280px', async ({ page }) => { + let verified = false + let resumed = false + let auditPage = 0 + const baseTask = { + id: 'task-a', account_id: 'account-a', draft_id: 'draft-a', confirmation_id: 'confirmation-a', state: 'needs_confirmation', + hold_reason: 'task_result_uncertain', browser_env_alias: 'env-a', network_exit_id: 'exit-a', runtime_instance_id: 'runtime-a', binding_version: 4, + confirmation: { id: 'confirmation-a', version: 1, account_version: 2, draft_version: 1, browser_env_alias: 'env-a', network_exit_id: 'exit-a', runtime_instance_id: 'runtime-a', binding_version: 4 }, + attempts: [{ id: 'attempt-a', started_at: '2026-08-30T00:00:00Z', outcome: 'uncertain', evidence: { mock_outcome: 'uncertain' } }], + } + const listTask = { ...baseTask, confirmation: undefined, attempts: undefined, allowed_action: undefined, updated_at: '2026-08-30T00:01:00Z' } + await page.route(/\/api\/phase-a\/tasks(?:\/.*)?(?:\?.*)?$/, async route => { + const request = route.request() + const url = new URL(request.url()) + if (request.method() === 'POST' && url.pathname.endsWith('/verify')) { + expect(request.postDataJSON()).toEqual({ result: 'not_executed' }) + verified = true + return route.fulfill({ status: 204 }) + } + if (request.method() === 'POST' && url.pathname.endsWith('/resume')) { + resumed = true + return route.fulfill({ status: 204 }) + } + if (url.pathname.endsWith('/task-a')) return route.fulfill({ json: { ...baseTask, state: resumed ? 'queued' : baseTask.state, allowed_action: resumed ? '' : verified ? 'resume' : 'verify', verification_result: verified ? 'not_executed' : undefined } }) + return route.fulfill({ json: [listTask] }) + }) + await page.route('**/api/phase-a/attempts/attempt-a', route => route.fulfill({ json: { ...baseTask.attempts[0], task_id: 'task-a', browser_env_alias: 'env-a', network_exit_id: 'exit-a', binding_version: 4 } })) + await page.route('**/api/browsers/env-a', route => route.fulfill({ json: { alias: 'env-a', name: '店铺环境', account_id: 'account-a', image_version: '148', binding_version: 4, runtime_instance_id: 'runtime-a', network_exit: { id: 'exit-a', health_status: 'healthy' } } })) + await page.route('**/api/network-exits/exit-a', route => route.fulfill({ json: { id: 'exit-a', protocol: 'socks5', host: 'proxy.example', port: 1080, health_status: 'healthy', version: 2 } })) + await page.route('**/api/phase-a/audit*', route => { + const currentPage = Number(new URL(route.request().url()).searchParams.get('page') || 1) + auditPage = Math.max(auditPage, currentPage) + return route.fulfill({ json: { data: [{ + id: 7, event_type: 'task_verified', account_id: 'account-a', task_id: 'task-a', attempt_id: 'attempt-a', browser_env_alias: 'env-a', network_exit_id: 'exit-a', binding_version: 4, reason_code: 'manual_verification_recorded', created_at: '2026-08-30T00:01:00Z', + }], total: 26, page: currentPage, page_size: 25 } }) + }) + + for (const path of ['/#/tasks', '/#/tasks/task-a', '/#/attempts/attempt-a', '/#/browsers/env-a', '/#/network-exits/exit-a', '/#/audit?task_id=task-a']) { + for (const width of [599, 900, 1280]) { + await page.setViewportSize({ width, height: 900 }) + await page.goto(path) + await expect(page.getByRole('heading', { level: 1 })).toBeVisible() + expect(await page.evaluate(() => document.documentElement.scrollWidth)).toBeLessThanOrEqual(width) + } + } + + verified = false + await page.setViewportSize({ width: 599, height: 900 }) + await page.goto('/#/tasks/task-a') + await expect(page.getByRole('button', { name: /恢复排队/ })).toHaveCount(0) + await page.getByRole('button', { name: '记录人工核验' }).click() + await expect(page.getByRole('button', { name: '按核验结论恢复排队' })).toBeVisible() + await page.getByRole('button', { name: '按核验结论恢复排队' }).click() + await expect.poll(() => resumed).toBe(true) + await expect(page.getByText('已排队')).toBeVisible() + await page.getByRole('link', { name: '查看关联审计' }).click() + await expect(page).toHaveURL(/#\/audit\?task_id=task-a$/) + await expect(page.getByRole('link', { name: 'task-a' }).first()).toBeVisible() + await page.getByRole('button', { name: '下一页' }).click() + await expect.poll(() => auditPage).toBe(2) + await page.getByRole('link', { name: 'attempt-a' }).first().click() + await expect(page).toHaveURL(/#\/attempts\/attempt-a$/) + await expect(page.getByRole('link', { name: '返回任务' })).toHaveAttribute('href', '#/tasks/task-a') +})