merge: allow account and competitor deletion
This commit is contained in:
@@ -0,0 +1,109 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/gofiber/fiber/v3"
|
||||
|
||||
"git.ipao.vip/rogee/creator-hub/internal/creator"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/hub"
|
||||
"git.ipao.vip/rogee/creator-hub/internal/phasea"
|
||||
)
|
||||
|
||||
func registerAccountDeletion(app *fiber.App, phaseAStore *phasea.Store, hubStore *hub.Store, creatorStore *creator.Store, credentials phasea.CredentialBridge) {
|
||||
app.Delete("/api/phase-a/accounts/:id", func(c fiber.Ctx) error {
|
||||
accountID := c.Params("id")
|
||||
if err := phaseAStore.CheckAccountDeletion(c.Context(), accountID); err != nil {
|
||||
return phaseAError(c, err)
|
||||
}
|
||||
if err := creatorStore.CheckOwnedAccountDeletion(c.Context(), accountID); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
|
||||
environment, hasEnvironment, err := accountEnvironment(c.Context(), hubStore, accountID)
|
||||
if err != nil {
|
||||
return hubError(c, err)
|
||||
}
|
||||
unlock := func() {}
|
||||
if hasEnvironment {
|
||||
aliases := []string{environment.Alias}
|
||||
exitIDs := []string{}
|
||||
if environment.Exit.ID != "" {
|
||||
exitIDs = append(exitIDs, environment.Exit.ID)
|
||||
}
|
||||
imageVersions := []string{}
|
||||
if environment.ImageVersion != "" {
|
||||
imageVersions = append(imageVersions, environment.ImageVersion)
|
||||
}
|
||||
unlock, err = hubStore.LockResources(c.Context(), aliases, exitIDs, imageVersions)
|
||||
if err != nil {
|
||||
return hubError(c, err)
|
||||
}
|
||||
defer unlock()
|
||||
gateway, gatewayErr := hubStore.GetGateway(c.Context(), environment.Gateway)
|
||||
if gatewayErr != nil {
|
||||
return hubError(c, gatewayErr)
|
||||
}
|
||||
if err := creatorStore.InvalidateListener(c.Context(), accountID, "账号删除"); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
if _, err := removeGatewayRuntime(c.Context(), hubStore, gateway, environment); err != nil {
|
||||
return hubError(c, err)
|
||||
}
|
||||
if err := purgeAccountProfile(c.Context(), gateway, environment); err != nil {
|
||||
return hubError(c, err)
|
||||
}
|
||||
}
|
||||
|
||||
if err := phaseAStore.DeleteAccountData(c.Context(), accountID); err != nil {
|
||||
return phaseAError(c, err)
|
||||
}
|
||||
if err := creatorStore.DeleteOwnedAccountData(c.Context(), accountID); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
if hasEnvironment {
|
||||
if err := hubStore.DeleteAccountEnvironment(c.Context(), accountID); err != nil {
|
||||
return hubError(c, err)
|
||||
}
|
||||
}
|
||||
if err := phaseAStore.DeleteAccount(c.Context(), accountID, credentials); err != nil {
|
||||
return phaseAError(c, err)
|
||||
}
|
||||
return c.SendStatus(fiber.StatusNoContent)
|
||||
})
|
||||
|
||||
app.Delete("/api/creator/competitors/:id", func(c fiber.Ctx) error {
|
||||
if err := creatorStore.DeleteCompetitor(c.Context(), c.Params("id")); err != nil {
|
||||
return creatorError(c, err)
|
||||
}
|
||||
return c.SendStatus(fiber.StatusNoContent)
|
||||
})
|
||||
}
|
||||
|
||||
func accountEnvironment(ctx context.Context, store *hub.Store, accountID string) (hub.EnvironmentContext, bool, error) {
|
||||
environment, err := store.GetEnvironmentContextForAccount(ctx, accountID)
|
||||
if errors.Is(err, hub.ErrNotFound) {
|
||||
return hub.EnvironmentContext{}, false, nil
|
||||
}
|
||||
if err != nil {
|
||||
return hub.EnvironmentContext{}, false, err
|
||||
}
|
||||
return environment, true, nil
|
||||
}
|
||||
|
||||
func purgeAccountProfile(ctx context.Context, gateway hub.Gateway, environment hub.EnvironmentContext) error {
|
||||
payload := gatewayCleanupGenerationPayload(environment)
|
||||
payload["purge_profile"] = true
|
||||
payload["profile_volume"] = "creatorhub-profile-" + environment.Alias
|
||||
status, body, err := gatewayCall(ctx, gateway, http.MethodDelete, "/v1/browsers/"+environment.Alias, payload, 30*time.Second)
|
||||
if err != nil {
|
||||
return gatewayUnreachable(err)
|
||||
}
|
||||
if status == http.StatusNoContent || status == http.StatusNotFound {
|
||||
return nil
|
||||
}
|
||||
return gatewayRejected(status, body)
|
||||
}
|
||||
@@ -327,6 +327,9 @@ func newHandlerWithCreatorAndAI(webDirectory, username, password string, phaseAS
|
||||
}
|
||||
registerCreatorWithServices(app, creatorStore, phaseAStore, hubStore, executor, generator, analyzer)
|
||||
}
|
||||
if phaseAStore != nil && hubStore != nil {
|
||||
registerAccountDeletion(app, phaseAStore, hubStore, creatorStore, credentials)
|
||||
}
|
||||
}
|
||||
app.Use(func(c fiber.Ctx) error {
|
||||
if isControlPlaneAPIPath(c.Path()) {
|
||||
|
||||
@@ -474,8 +474,17 @@ class Gateway:
|
||||
|
||||
def remove(self, alias: str, input: dict) -> None:
|
||||
generation = decode_generation(
|
||||
input, require_runtime=False, require_network=False
|
||||
input,
|
||||
require_runtime=False,
|
||||
require_network=False,
|
||||
allow_profile_purge=True,
|
||||
)
|
||||
purge_profile = input.get("purge_profile", False)
|
||||
profile_volume = input.get("profile_volume", "")
|
||||
if type(purge_profile) is not bool or not isinstance(profile_volume, str):
|
||||
raise RequestError("profile purge fields are invalid", 400)
|
||||
if purge_profile and profile_volume != f"creatorhub-profile-{alias}":
|
||||
raise RequestError("profile volume does not belong to browser alias", 400)
|
||||
with self._alias_lock(alias):
|
||||
try:
|
||||
container_id, labels = self.docker.managed_container(alias)
|
||||
@@ -543,6 +552,16 @@ class Gateway:
|
||||
pass
|
||||
except Exception as exc:
|
||||
raise RequestError("Docker container removal failed") from exc
|
||||
if purge_profile:
|
||||
try:
|
||||
self.docker.expect(
|
||||
"DELETE",
|
||||
f"/volumes/{quote(profile_volume, safe='')}",
|
||||
)
|
||||
except FileNotFoundError:
|
||||
pass
|
||||
except Exception as exc:
|
||||
raise RequestError("Docker profile volume removal failed") from exc
|
||||
self._release_action_ownership(alias)
|
||||
|
||||
def restore_proxy(self, alias: str, input: dict) -> None:
|
||||
@@ -1657,9 +1676,14 @@ def has_control(value: str) -> bool:
|
||||
|
||||
|
||||
def decode_generation(
|
||||
value: dict, require_runtime: bool, require_network: bool
|
||||
value: dict,
|
||||
require_runtime: bool,
|
||||
require_network: bool,
|
||||
allow_profile_purge: bool = False,
|
||||
) -> dict:
|
||||
allowed = {"binding_version", "runtime_id", "network_id"}
|
||||
if allow_profile_purge:
|
||||
allowed.update({"purge_profile", "profile_volume"})
|
||||
if not isinstance(value, dict) or set(value) - allowed:
|
||||
raise RequestError(
|
||||
"binding_version, runtime_id and network_id must identify the expected generation",
|
||||
|
||||
+3
-3
@@ -96,7 +96,7 @@ docker compose exec -T postgres \
|
||||
| grep -qx 1
|
||||
docker compose exec -T postgres \
|
||||
psql -U creatorhub -d creatorhub -tAc \
|
||||
'SELECT 1 FROM schema_migration WHERE version = 33;' \
|
||||
'SELECT 1 FROM schema_migration WHERE version = 35;' \
|
||||
| grep -qx 1
|
||||
|
||||
docker compose ps
|
||||
@@ -419,7 +419,7 @@ docker compose up --detach --build
|
||||
| grep -qx 1
|
||||
docker compose exec -T postgres \
|
||||
psql -U creatorhub -d creatorhub_restore_check -v ON_ERROR_STOP=1 -tAc \
|
||||
'SELECT 1 FROM schema_migration WHERE version = 33;' \
|
||||
'SELECT 1 FROM schema_migration WHERE version = 35;' \
|
||||
| grep -qx 1
|
||||
docker compose exec -T postgres \
|
||||
dropdb --force -U creatorhub creatorhub_restore_check
|
||||
@@ -470,7 +470,7 @@ docker compose up --detach --build
|
||||
| grep -qx 1
|
||||
docker compose exec -T postgres \
|
||||
psql -U creatorhub -d creatorhub -v ON_ERROR_STOP=1 -tAc \
|
||||
'SELECT 1 FROM schema_migration WHERE version = 33;' \
|
||||
'SELECT 1 FROM schema_migration WHERE version = 35;' \
|
||||
| grep -qx 1
|
||||
|
||||
git switch --detach "$RESTORE_REV"
|
||||
|
||||
@@ -0,0 +1,288 @@
|
||||
package creator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"errors"
|
||||
"time"
|
||||
)
|
||||
|
||||
func (s *Store) CheckOwnedAccountDeletion(ctx context.Context, accountID string) error {
|
||||
if accountID == "" {
|
||||
return ErrInvalid
|
||||
}
|
||||
var exists bool
|
||||
if err := s.db.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM social_account WHERE id = $1)`, accountID).Scan(&exists); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if !exists {
|
||||
return ErrNotFound
|
||||
}
|
||||
var active bool
|
||||
if err := s.db.QueryRowContext(ctx, `
|
||||
SELECT EXISTS (
|
||||
SELECT 1 FROM creator_operation
|
||||
WHERE account_id = $1 AND state = 'processing'
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_event
|
||||
WHERE (receiving_account_id = $1 OR execution_account_id = $1) AND state = 'processing'
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_collection_checkpoint
|
||||
WHERE source_type = 'owned' AND source_id = $1 AND status = 'running'
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_source_sync_lease
|
||||
WHERE source_type = 'owned' AND source_id = $1 AND lease_until > now()
|
||||
)`, accountID).Scan(&active); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if active {
|
||||
return ErrConflict
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Store) DeleteOwnedAccountData(ctx context.Context, accountID string) error {
|
||||
if accountID == "" {
|
||||
return ErrInvalid
|
||||
}
|
||||
tx, err := s.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return errors.New("begin creator account deletion")
|
||||
}
|
||||
defer tx.Rollback()
|
||||
if err := lockCreatorAccount(ctx, tx, accountID); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := checkOwnedAccountDeletionTx(ctx, tx, accountID); err != nil {
|
||||
return err
|
||||
}
|
||||
var secretID, secretKey string
|
||||
var hasSecret bool
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT secret_reference_id, secret_key
|
||||
FROM creator_account_password
|
||||
WHERE account_id = $1`, accountID).Scan(&secretID, &secretKey); err == nil {
|
||||
hasSecret = true
|
||||
} else if !errors.Is(err, sql.ErrNoRows) {
|
||||
return databaseError(err)
|
||||
}
|
||||
if hasSecret && s.secrets == nil {
|
||||
return errors.New("creator secret bridge is unavailable")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_operation
|
||||
WHERE account_id = $1
|
||||
OR strategy_id IN (SELECT id FROM creator_strategy WHERE big_account_id = $1 OR execution_account_id = $1)
|
||||
OR event_id IN (SELECT id FROM creator_event WHERE receiving_account_id = $1 OR execution_account_id = $1)`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_event
|
||||
WHERE receiving_account_id = $1 OR execution_account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_cooldown
|
||||
WHERE big_account_id = $1 OR execution_account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_strategy
|
||||
WHERE big_account_id = $1 OR execution_account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_relation
|
||||
WHERE big_account_id = $1 OR small_account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_conversation WHERE account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_listener_boundary WHERE account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_listener_state WHERE account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if err := deleteCreatorSourceTx(ctx, tx, SourceOwned, accountID); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_account_password WHERE account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_account_profile WHERE account_id = $1`, accountID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return errors.New("commit creator account deletion")
|
||||
}
|
||||
if hasSecret {
|
||||
if err := s.secrets.Delete(context.WithoutCancel(ctx), SecretReference{ID: secretID, Provider: "os_keyring"}, secretKey); err != nil {
|
||||
return errors.Join(errors.New("delete creator account password"), err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Store) DeleteCompetitor(ctx context.Context, competitorID string) error {
|
||||
if competitorID == "" {
|
||||
return ErrInvalid
|
||||
}
|
||||
tx, err := s.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return errors.New("begin competitor deletion")
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var syncStatus string
|
||||
var syncLease sql.NullTime
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT sync_status, sync_lease_until
|
||||
FROM creator_competitor
|
||||
WHERE id = $1
|
||||
FOR UPDATE`, competitorID).Scan(&syncStatus, &syncLease); err != nil {
|
||||
return rowError(err)
|
||||
}
|
||||
if syncStatus == "running" || (syncLease.Valid && syncLease.Time.After(time.Now())) {
|
||||
return ErrConflict
|
||||
}
|
||||
var active bool
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT EXISTS (
|
||||
SELECT 1 FROM creator_collection_checkpoint
|
||||
WHERE source_type = 'competitor' AND source_id = $1
|
||||
AND (status = 'running' OR (lease_until IS NOT NULL AND lease_until > now()))
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_source_sync_lease
|
||||
WHERE source_type = 'competitor' AND source_id = $1 AND lease_until > now()
|
||||
)`, competitorID).Scan(&active); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if active {
|
||||
return ErrConflict
|
||||
}
|
||||
if err := deleteCreatorSourceTx(ctx, tx, SourceCompetitor, competitorID); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_competitor WHERE id = $1`, competitorID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return errors.New("commit competitor deletion")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func lockCreatorAccount(ctx context.Context, tx *sql.Tx, accountID string) error {
|
||||
var lockedID string
|
||||
if err := tx.QueryRowContext(ctx, `SELECT id FROM social_account WHERE id = $1 FOR UPDATE`, accountID).Scan(&lockedID); err != nil {
|
||||
return rowError(err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func checkOwnedAccountDeletionTx(ctx context.Context, tx *sql.Tx, accountID string) error {
|
||||
var active bool
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT EXISTS (
|
||||
SELECT 1 FROM creator_operation
|
||||
WHERE account_id = $1 AND state = 'processing'
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_event
|
||||
WHERE (receiving_account_id = $1 OR execution_account_id = $1) AND state = 'processing'
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_collection_checkpoint
|
||||
WHERE source_type = 'owned' AND source_id = $1 AND status = 'running'
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_source_sync_lease
|
||||
WHERE source_type = 'owned' AND source_id = $1 AND lease_until > now()
|
||||
)`, accountID).Scan(&active); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if active {
|
||||
return ErrConflict
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func deleteCreatorSourceTx(ctx context.Context, tx *sql.Tx, sourceType, sourceID string) error {
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
CREATE TEMP TABLE creator_delete_work_ids (id text PRIMARY KEY) ON COMMIT DROP`); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
INSERT INTO creator_delete_work_ids (id)
|
||||
SELECT id FROM creator_work WHERE source_type = $1 AND source_id = $2
|
||||
UNION
|
||||
SELECT work_id FROM creator_work_source WHERE source_type = $1 AND source_id = $2`, sourceType, sourceID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
var active bool
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT EXISTS (
|
||||
SELECT 1 FROM creator_operation
|
||||
WHERE state = 'processing'
|
||||
AND (target_work_id IN (SELECT id FROM creator_delete_work_ids)
|
||||
OR target_comment_id IN (SELECT id FROM creator_comment WHERE work_id IN (SELECT id FROM creator_delete_work_ids)))
|
||||
UNION ALL
|
||||
SELECT 1 FROM creator_event
|
||||
WHERE state = 'processing' AND work_id IN (SELECT id FROM creator_delete_work_ids)
|
||||
)`).Scan(&active); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if active {
|
||||
return ErrConflict
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_work_source WHERE source_type = $1 AND source_id = $2`, sourceType, sourceID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_source_sync_lease WHERE source_type = $1 AND source_id = $2`, sourceType, sourceID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM creator_collection_checkpoint WHERE source_type = $1 AND source_id = $2`, sourceType, sourceID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
UPDATE creator_work work
|
||||
SET source_type = source.source_type, source_id = source.source_id
|
||||
FROM (
|
||||
SELECT DISTINCT ON (work_id) work_id, source_type, source_id
|
||||
FROM creator_work_source
|
||||
ORDER BY work_id, created_at, source_type, source_id
|
||||
) source
|
||||
WHERE work.id = source.work_id AND work.source_type = $1 AND work.source_id = $2`, sourceType, sourceID); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_operation
|
||||
WHERE target_work_id IN (
|
||||
SELECT id FROM creator_delete_work_ids
|
||||
WHERE NOT EXISTS (SELECT 1 FROM creator_work_source WHERE work_id = creator_delete_work_ids.id)
|
||||
)
|
||||
OR target_comment_id IN (
|
||||
SELECT comment.id FROM creator_comment comment
|
||||
WHERE comment.work_id IN (
|
||||
SELECT id FROM creator_delete_work_ids
|
||||
WHERE NOT EXISTS (SELECT 1 FROM creator_work_source WHERE work_id = creator_delete_work_ids.id)
|
||||
)
|
||||
)`); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_event
|
||||
WHERE work_id IN (
|
||||
SELECT id FROM creator_delete_work_ids
|
||||
WHERE NOT EXISTS (SELECT 1 FROM creator_work_source WHERE work_id = creator_delete_work_ids.id)
|
||||
)`); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM creator_work
|
||||
WHERE id IN (
|
||||
SELECT id FROM creator_delete_work_ids
|
||||
WHERE NOT EXISTS (SELECT 1 FROM creator_work_source WHERE work_id = creator_delete_work_ids.id)
|
||||
)`); err != nil {
|
||||
return databaseError(err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
ALTER TABLE creator_competitor
|
||||
ADD COLUMN IF NOT EXISTS tags text[] NOT NULL DEFAULT ARRAY[]::text[];
|
||||
|
||||
ALTER TABLE creator_competitor
|
||||
DROP CONSTRAINT IF EXISTS creator_competitor_tags_check,
|
||||
ADD CONSTRAINT creator_competitor_tags_check
|
||||
CHECK (cardinality(tags) <= 20);
|
||||
@@ -0,0 +1,9 @@
|
||||
CREATE OR REPLACE FUNCTION reject_audit_event_mutation() RETURNS trigger
|
||||
LANGUAGE plpgsql AS $$
|
||||
BEGIN
|
||||
IF current_setting('creatorhub.account_deletion', true) = 'on' THEN
|
||||
RETURN OLD;
|
||||
END IF;
|
||||
RAISE EXCEPTION 'audit_event is append-only';
|
||||
END;
|
||||
$$;
|
||||
@@ -67,6 +67,16 @@ var migration032 string
|
||||
//go:embed migrations/033_competitor_tags.sql
|
||||
var migration033 string
|
||||
|
||||
// schema_migration is shared by phase-a, hub, and creator. Hub already used
|
||||
// version 33, so version 34 repairs creator competitor tags when creator
|
||||
// migration 33 was skipped.
|
||||
//
|
||||
//go:embed migrations/034_competitor_tags_repair.sql
|
||||
var migration034 string
|
||||
|
||||
//go:embed migrations/035_account_deletion.sql
|
||||
var migration035 string
|
||||
|
||||
type SecretReference struct {
|
||||
ID string
|
||||
Provider string
|
||||
@@ -161,6 +171,8 @@ func (s *Store) migrate(ctx context.Context) error {
|
||||
{version: 31, sql: migration031},
|
||||
{version: 32, sql: migration032},
|
||||
{version: 33, sql: migration033},
|
||||
{version: 34, sql: migration034},
|
||||
{version: 35, sql: migration035},
|
||||
}
|
||||
for _, migration := range migrations {
|
||||
var applied bool
|
||||
|
||||
@@ -548,6 +548,36 @@ func (s *Store) DeleteEnv(ctx context.Context, alias string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteAccountEnvironment removes the account binding and its private browser environment.
|
||||
// Shared gateways, images, and network exits are intentionally preserved.
|
||||
func (s *Store) DeleteAccountEnvironment(ctx context.Context, accountID string) error {
|
||||
if strings.TrimSpace(accountID) == "" {
|
||||
return ErrInvalid
|
||||
}
|
||||
tx, err := s.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return errors.New("begin account environment deletion")
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var alias string
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT browser_env_alias
|
||||
FROM environment_binding
|
||||
WHERE account_id = $1
|
||||
FOR UPDATE`, accountID).Scan(&alias); errors.Is(err, sql.ErrNoRows) {
|
||||
return nil
|
||||
} else if err != nil {
|
||||
return errors.New("read account environment binding")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM environment_binding WHERE account_id = $1`, accountID); err != nil {
|
||||
return errors.New("delete account environment binding")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM browser_env WHERE alias = $1`, alias); err != nil {
|
||||
return errors.New("delete account browser environment")
|
||||
}
|
||||
return commitHub(tx)
|
||||
}
|
||||
|
||||
func scanEnv(rows *sql.Rows) (Env, error) {
|
||||
var env Env
|
||||
var encoded []byte
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
package phasea
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
)
|
||||
|
||||
// CheckAccountDeletion verifies that deleting an account will not interrupt an active run.
|
||||
func (s *Store) CheckAccountDeletion(ctx context.Context, accountID string) error {
|
||||
if !idPattern.MatchString(accountID) {
|
||||
return ErrInvalid
|
||||
}
|
||||
var exists bool
|
||||
if err := s.db.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM social_account WHERE id = $1)`, accountID).Scan(&exists); err != nil {
|
||||
return errors.New("check account deletion state")
|
||||
}
|
||||
if !exists {
|
||||
return ErrNotFound
|
||||
}
|
||||
var active bool
|
||||
if err := s.db.QueryRowContext(ctx, `
|
||||
SELECT EXISTS (
|
||||
SELECT 1 FROM runtime_instance
|
||||
WHERE account_id = $1 AND released_at IS NULL
|
||||
UNION ALL
|
||||
SELECT 1 FROM operation_task
|
||||
WHERE account_id = $1 AND state = 'executing'
|
||||
)`, accountID).Scan(&active); err != nil {
|
||||
return errors.New("check account deletion state")
|
||||
}
|
||||
if active {
|
||||
return ErrConflict
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteAccountData removes Phase A data while keeping the account row for the final deletion step.
|
||||
func (s *Store) DeleteAccountData(ctx context.Context, accountID string) error {
|
||||
if !idPattern.MatchString(accountID) {
|
||||
return ErrInvalid
|
||||
}
|
||||
tx, err := s.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return errors.New("begin account deletion transaction")
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var lockedID string
|
||||
if err := tx.QueryRowContext(ctx, `SELECT id FROM social_account WHERE id = $1 FOR UPDATE`, accountID).Scan(&lockedID); err != nil {
|
||||
return rowError(err)
|
||||
}
|
||||
var active bool
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT EXISTS (
|
||||
SELECT 1 FROM runtime_instance
|
||||
WHERE account_id = $1 AND released_at IS NULL
|
||||
UNION ALL
|
||||
SELECT 1 FROM operation_task
|
||||
WHERE account_id = $1 AND state = 'executing'
|
||||
)`, accountID).Scan(&active); err != nil {
|
||||
return errors.New("check account deletion state")
|
||||
}
|
||||
if active {
|
||||
return ErrConflict
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `SELECT set_config('creatorhub.account_deletion', 'on', true)`); err != nil {
|
||||
return errors.New("enable account deletion audit cleanup")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM audit_event
|
||||
WHERE account_id = $1
|
||||
OR confirmation_id IN (SELECT id FROM confirmation WHERE account_id = $1)
|
||||
OR task_id IN (SELECT id FROM operation_task WHERE account_id = $1)
|
||||
OR attempt_id IN (
|
||||
SELECT id FROM execution_attempt
|
||||
WHERE task_id IN (SELECT id FROM operation_task WHERE account_id = $1)
|
||||
)
|
||||
OR browser_env_alias IN (
|
||||
SELECT browser_env_alias FROM environment_binding WHERE account_id = $1
|
||||
)
|
||||
OR runtime_instance_id IN (SELECT id FROM runtime_instance WHERE account_id = $1)`, accountID); err != nil {
|
||||
return errors.New("delete account audit data")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `UPDATE operation_task SET current_attempt_id = NULL WHERE account_id = $1`, accountID); err != nil {
|
||||
return errors.New("detach account task attempts")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `
|
||||
DELETE FROM execution_attempt
|
||||
WHERE task_id IN (SELECT id FROM operation_task WHERE account_id = $1)`, accountID); err != nil {
|
||||
return errors.New("delete account task attempts")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM operation_task WHERE account_id = $1`, accountID); err != nil {
|
||||
return errors.New("delete account tasks")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM confirmation WHERE account_id = $1`, accountID); err != nil {
|
||||
return errors.New("delete account confirmations")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM content_draft WHERE account_id = $1`, accountID); err != nil {
|
||||
return errors.New("delete account drafts")
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM runtime_instance WHERE account_id = $1`, accountID); err != nil {
|
||||
return errors.New("delete account runtime records")
|
||||
}
|
||||
return commit(tx)
|
||||
}
|
||||
|
||||
// DeleteAccount removes the account row and its external cookie credential.
|
||||
// Call DeleteAccountData and delete account-owned creator/environment data first.
|
||||
func (s *Store) DeleteAccount(ctx context.Context, accountID string, credentials CredentialBridge) error {
|
||||
if !idPattern.MatchString(accountID) || credentials == nil {
|
||||
return ErrInvalid
|
||||
}
|
||||
tx, err := s.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return errors.New("begin account record deletion")
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var reference CredentialReference
|
||||
var key string
|
||||
if err := tx.QueryRowContext(ctx, `
|
||||
SELECT credential.id, credential.provider, credential.reference_key
|
||||
FROM social_account account
|
||||
JOIN credential_reference credential ON credential.id = account.credential_reference_id
|
||||
WHERE account.id = $1
|
||||
FOR UPDATE`, accountID).Scan(&reference.ID, &reference.Provider, &key); err != nil {
|
||||
return rowError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM social_account WHERE id = $1`, accountID); err != nil {
|
||||
return publicDatabaseError(err)
|
||||
}
|
||||
if _, err := tx.ExecContext(ctx, `DELETE FROM credential_reference WHERE id = $1`, reference.ID); err != nil {
|
||||
return publicDatabaseError(err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return errors.New("commit account record deletion")
|
||||
}
|
||||
if err := credentials.Delete(context.WithoutCancel(ctx), reference, key); err != nil {
|
||||
return errors.Join(errors.New("delete account credential"), err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -323,6 +323,7 @@ export function AccountList() {
|
||||
const [tagTarget, setTagTarget] = useState(null);
|
||||
const [tagDraft, setTagDraft] = useState([]);
|
||||
const [tagBusy, setTagBusy] = useState(false);
|
||||
const [deleteBusy, setDeleteBusy] = useState("");
|
||||
useTitle("CreatorHub · 账号管理");
|
||||
|
||||
const load = async () => {
|
||||
@@ -419,6 +420,31 @@ export function AccountList() {
|
||||
}
|
||||
}
|
||||
|
||||
async function deleteAccount(row) {
|
||||
const label = row.source_type === "owned" ? "自有账号" : "监测账号";
|
||||
if (!window.confirm(`确定删除${label}“${row.name}”吗?相关专属资源也会被清理。`)) {
|
||||
return;
|
||||
}
|
||||
const busyID = `${row.source_type}:${row.id}`;
|
||||
setDeleteBusy(busyID);
|
||||
setNotice(null);
|
||||
try {
|
||||
await dataProvider.deleteOne({
|
||||
resource: row.source_type === "owned" ? "accounts" : "creator-competitors",
|
||||
id: row.id,
|
||||
});
|
||||
setNotice({ variant: "success", text: `${label}已删除。` });
|
||||
await load();
|
||||
} catch (deleteError) {
|
||||
setNotice({
|
||||
variant: "destructive",
|
||||
text: conflictMessage(deleteError, `${label}删除失败`),
|
||||
});
|
||||
} finally {
|
||||
setDeleteBusy("");
|
||||
}
|
||||
}
|
||||
|
||||
function openTags(row) {
|
||||
setNotice(null);
|
||||
setTagTarget(row);
|
||||
@@ -638,6 +664,13 @@ export function AccountList() {
|
||||
</Button>
|
||||
</>
|
||||
)}
|
||||
<Button
|
||||
size="sm"
|
||||
onClick={() => deleteAccount(account)}
|
||||
busy={deleteBusy === `${account.source_type}:${account.id}`}
|
||||
>
|
||||
删除
|
||||
</Button>
|
||||
</div>
|
||||
),
|
||||
},
|
||||
|
||||
@@ -65,6 +65,7 @@ function provider(overrides = {}) {
|
||||
update: vi.fn(),
|
||||
updateMany: vi.fn(),
|
||||
delete: vi.fn(),
|
||||
deleteOne: vi.fn().mockResolvedValue({ data: {} }),
|
||||
deleteMany: vi.fn(),
|
||||
...overrides,
|
||||
};
|
||||
@@ -246,6 +247,46 @@ describe("AccountList", () => {
|
||||
),
|
||||
).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("deletes owned and monitoring accounts through their own resources", async () => {
|
||||
const competitor = {
|
||||
id: "competitor-a",
|
||||
nickname: "竞品账号",
|
||||
platform: "douyin",
|
||||
platform_account_key: "competitor-key",
|
||||
tags: [],
|
||||
enabled: true,
|
||||
};
|
||||
const dataProvider = provider({
|
||||
getList: vi.fn(({ resource }) =>
|
||||
resource === "accounts"
|
||||
? { data: [account], total: 1 }
|
||||
: resource === "creator-competitors"
|
||||
? { data: [competitor], total: 1 }
|
||||
: { data: [], total: 0 },
|
||||
),
|
||||
deleteOne: vi.fn().mockResolvedValue({ data: {} }),
|
||||
});
|
||||
vi.spyOn(window, "confirm").mockReturnValue(true);
|
||||
renderList(dataProvider);
|
||||
|
||||
expect(await screen.findByText("竞品账号")).toBeTruthy();
|
||||
const deleteButtons = screen.getAllByRole("button", { name: "删除" });
|
||||
fireEvent.click(deleteButtons[0]);
|
||||
await waitFor(() =>
|
||||
expect(dataProvider.deleteOne).toHaveBeenCalledWith({
|
||||
resource: "accounts",
|
||||
id: "account-a",
|
||||
}),
|
||||
);
|
||||
fireEvent.click(screen.getAllByRole("button", { name: "删除" })[1]);
|
||||
await waitFor(() =>
|
||||
expect(dataProvider.deleteOne).toHaveBeenCalledWith({
|
||||
resource: "creator-competitors",
|
||||
id: "competitor-a",
|
||||
}),
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
describe("AccountDetail", () => {
|
||||
|
||||
@@ -162,7 +162,10 @@ export const dataProvider = {
|
||||
},
|
||||
async deleteOne({ resource, id }) {
|
||||
const path = resourcePaths[resource];
|
||||
if (!path || resource !== "browser-images")
|
||||
if (
|
||||
!path ||
|
||||
!["browser-images", "accounts", "creator-competitors"].includes(resource)
|
||||
)
|
||||
return unsupported(resource, "deleteOne");
|
||||
await request(`${path}/${encodeURIComponent(id)}`, { method: "DELETE" });
|
||||
return { data: { id } };
|
||||
|
||||
Reference in New Issue
Block a user