import argparse import json import sys from pathlib import Path from urllib.error import HTTPError, URLError from urllib.request import Request, urlopen import psycopg from psycopg.rows import dict_row from wiki_update_page import DEFAULT_WIKI_URL, WikiUpdateError, connect_info, find_auth_group, load_env, login TEMP_PERMISSIONS = ["manage:system", "manage:assets", "read:assets", "write:assets"] def normalize_json(value): if value is None or isinstance(value, (dict, list)): return value return json.loads(value) def get_json_cast(cur, column_name): cur.execute( """ SELECT udt_name FROM information_schema.columns WHERE table_name = 'groups' AND column_name = %s """, (column_name,), ) row = cur.fetchone() if not row: raise WikiUpdateError(f"groups.{column_name} was not found") return "jsonb" if row["udt_name"] == "jsonb" else "json" def grant_temporary_permissions(cur, group): permissions_cast = get_json_cast(cur, "permissions") permissions = normalize_json(group["permissions"]) or [] added = [] for permission in TEMP_PERMISSIONS: if permission not in permissions: permissions.append(permission) added.append(permission) if added: cur.execute( f"UPDATE groups SET permissions = %s::{permissions_cast} WHERE id = %s", (json.dumps(permissions), group["id"]), ) print(f"temporary_permissions_added={','.join(added) if added else 'none'}") return added def remove_temporary_permissions(conninfo, group_id, added): if not added: return with psycopg.connect(**conninfo, row_factory=dict_row) as conn: with conn.cursor() as cur: permissions_cast = get_json_cast(cur, "permissions") cur.execute("SELECT permissions FROM groups WHERE id = %s", (group_id,)) row = cur.fetchone() permissions = normalize_json(row["permissions"]) or [] permissions = [permission for permission in permissions if permission not in added] cur.execute( f"UPDATE groups SET permissions = %s::{permissions_cast} WHERE id = %s", (json.dumps(permissions), group_id), ) conn.commit() print(f"temporary_permissions_removed={','.join(added)}") def misplaced_assets(cur): cur.execute( """ SELECT a.id, a.filename, a.hash, a."fileSize", f.slug AS folder_slug FROM assets a JOIN "assetFolders" f ON f.id = a."folderId" WHERE f.slug = 'vulture' AND (left(a.filename, 5) = 'echo_' OR left(a.filename, 11) = 'enterprise_') ORDER BY a.filename, a.id """ ) return cur.fetchall() def graphql(wiki_url, query, variables, token): body = json.dumps({"query": query, "variables": variables}).encode("utf-8") request = Request( f"{wiki_url}/graphql", data=body, headers={"Content-Type": "application/json", "Accept": "application/json", "Authorization": f"Bearer {token}"}, method="POST", ) try: with urlopen(request, timeout=120) as response: return json.loads(response.read().decode("utf-8")) except HTTPError as exc: return {"transportError": f"HTTP {exc.code}", "payload": exc.read().decode("utf-8", errors="replace")} except URLError as exc: return {"transportError": str(exc)} def delete_asset(wiki_url, token, asset_id): mutation = """ mutation DeleteAsset($id: Int!) { assets { deleteAsset(id: $id) { responseResult { succeeded errorCode slug message } } } } """ payload = graphql(wiki_url, mutation, {"id": asset_id}, token) if payload.get("errors"): print(f"asset_delete id={asset_id} graphql_errors={json.dumps(payload['errors'], ensure_ascii=True)}") return False if payload.get("transportError"): print(f"asset_delete id={asset_id} transport_error={payload['transportError']} payload={payload.get('payload', '')[:500]}") return False result = payload.get("data", {}).get("assets", {}).get("deleteAsset", {}).get("responseResult", {}) print( "asset_delete id={id} succeeded={succeeded} errorCode={errorCode} slug={slug} message={message}".format( id=asset_id, succeeded=result.get("succeeded"), errorCode=result.get("errorCode"), slug=result.get("slug"), message=result.get("message"), ) ) return bool(result.get("succeeded")) def parse_args(): parser = argparse.ArgumentParser(description="Delete Echo/Enterprise assets that were mistakenly uploaded under the Vulture folder.") parser.add_argument("--env-file", default=".env") parser.add_argument("--wiki-url", default=DEFAULT_WIKI_URL) parser.add_argument("--auth-group", default="automation") parser.add_argument("--delete", action="store_true", help="Actually delete matching assets. Without this flag, only lists matches.") return parser.parse_args() def main(): args = parse_args() repo_root = Path.cwd() env = load_env((repo_root / args.env_file).resolve()) conninfo = connect_info(env) identities = [] for candidate in (env.get("wiki_useremail"), env.get("wiki_username")): if candidate and candidate not in identities: identities.append(candidate) if not identities or not env.get("wiki_password"): raise WikiUpdateError("Missing wiki_useremail/wiki_username or wiki_password in .env") temporary_group_id = None added_permissions = [] try: with psycopg.connect(**conninfo, row_factory=dict_row) as conn: with conn.cursor() as cur: assets = misplaced_assets(cur) print(f"misplaced_asset_count={len(assets)}") for asset in assets: print(f"misplaced_asset id={asset['id']} folder=/{asset['folder_slug']} filename={asset['filename']} size={asset['fileSize']}") if not args.delete or not assets: return 0 group = find_auth_group(cur, args.auth_group, [identity.lower() for identity in identities]) added_permissions = grant_temporary_permissions(cur, group) temporary_group_id = group["id"] conn.commit() token = login(args.wiki_url.rstrip("/"), identities, env["wiki_password"]) ok = True for asset in assets: ok = delete_asset(args.wiki_url.rstrip("/"), token, asset["id"]) and ok if not ok: return 1 with psycopg.connect(**conninfo, row_factory=dict_row) as conn: with conn.cursor() as cur: remaining = misplaced_assets(cur) print(f"misplaced_asset_remaining={len(remaining)}") return 0 if not remaining else 1 finally: if temporary_group_id is not None: remove_temporary_permissions(conninfo, temporary_group_id, added_permissions) if __name__ == "__main__": try: sys.exit(main()) except WikiUpdateError as exc: print(f"wiki_cleanup_error={exc}") sys.exit(2)