diff --git a/.gitignore b/.gitignore index 4fd661b..d6c8aa4 100644 --- a/.gitignore +++ b/.gitignore @@ -2,6 +2,8 @@ node_modules/ web/node_modules/ web/dist/ +web/.output/ +web/src/routeTree.gen.ts web/.umi/ web/.umi-production/ web/.umi-test/ @@ -9,6 +11,14 @@ web/.src/ dist/ build/ +# 测试报告与缓存 +web/coverage/ +web/test-results/ +.pytest_cache/ +.ruff_cache/ +.vitest/ +web/.tanstack/ + # 日志与临时文件 *.log tmp/ diff --git a/AGENTS.md b/AGENTS.md index c987f47..584dccd 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -39,6 +39,8 @@ 产品方向:平台当前仅支持抖音(douyin);小红书等其他平台的业务与代码已全部移除,不得重新引入。 +环境先行:创建时只创建抖音浏览器环境,不填写昵称、UID 或 Cookie,不创建占位账号。未登录环境在账号列表独立显示。首次浏览器身份核验成功后,同一事务创建/关联账号并同步真实 UID、昵称、头像、抖音号与 secUID;同一 UID 只能绑定一个账号,已有绑定禁止换绑或自动覆盖。后续同 UID 核验更新平台资料,不覆盖本地备注和业务设置。浏览器 profile_id 与指纹 seed 在环境创建时确定,账号绑定不得改变它们;旧环境保持原 profile_id 和 seed,不重建浏览器或清空 Cookie。指纹 seed 由环境独立序列分配,不再依赖账号 ID。 + 前端框架:Umi Max 4.7 + React 19; 组件库:antd 6.6.5 + @ant-design/pro-components 3.x(beta 线)+ @ant-design/icons;仅使用 antd/pro 默认组件原样实现,禁止自定义封装与样式魔改;组件不满足业务时改交互逻辑适配组件; diff --git a/browser_gateway/browser/cdp.py b/browser_gateway/browser/cdp.py index 59d73d4..36220e3 100644 --- a/browser_gateway/browser/cdp.py +++ b/browser_gateway/browser/cdp.py @@ -68,7 +68,7 @@ class CDPConnection: self._pending.append(message) continue if message.get("error"): - raise BrowserError(f"CDP command rejected: {method}") + raise BrowserError(f"CDP command rejected: {method}: {json.dumps(message['error'], ensure_ascii=False)}") result = message.get("result") if not isinstance(result, dict): raise BrowserError(f"CDP command returned invalid result: {method}") diff --git a/browser_gateway/platform/douyin.py b/browser_gateway/platform/douyin.py index 52a90a5..0134df5 100644 --- a/browser_gateway/platform/douyin.py +++ b/browser_gateway/platform/douyin.py @@ -85,6 +85,11 @@ ACTIONS = frozenset( DouyinError = BrowserError + +class DouyinLoginPending(BrowserError): + """The browser is healthy but the user has not completed login yet.""" + + LISTENER_ERRORS = ( DouyinError, OSError, @@ -364,6 +369,9 @@ class DouyinBrowser: return final_url def _wait_for_login_render(self, cdp: CDPConnection) -> None: + # Page.navigate acknowledges navigation before a document/view exists. + # Wait for the new DOM before requesting its first screenshot. + cdp.wait_event("Page.domContentEventFired", lambda params: True, timeout=LOGIN_RENDER_TIMEOUT) deadline = time.monotonic() + LOGIN_RENDER_TIMEOUT while True: screenshot = cdp.command( @@ -381,6 +389,7 @@ class DouyinBrowser: def login_qr(self, alias: str) -> BrowserLoginQRResponse: with self.connection(alias) as cdp: + cdp.command("Page.enable") navigation = cdp.command("Page.navigate", {"url": LOGIN_PAGE_URL}) if ( not isinstance(navigation, dict) @@ -401,8 +410,7 @@ class DouyinBrowser: ) if isinstance(opened, bool) and opened: time.sleep(0.5) - page = cdp.evaluate( - """(() => { + qr_expression = """(() => { const visible = element => { const rect = element.getBoundingClientRect(); const style = getComputedStyle(element); @@ -411,6 +419,7 @@ class DouyinBrowser: }; const qrElement = [...document.querySelectorAll('img,canvas')].find(element => { if (!visible(element)) return false; + if (element.tagName === 'IMG' && (!element.complete || element.naturalWidth === 0)) return false; const rect = element.getBoundingClientRect(); const label = `${element.alt || ''} ${element.title || ''} ${element.getAttribute('aria-label') || ''}`; return /二维码|qr.?code/i.test(label) || @@ -432,16 +441,23 @@ class DouyinBrowser: }, }; })()""" - ) - if ( - not isinstance(page, dict) - or page.get("origin") not in LOGIN_ORIGINS - or not isinstance(page.get("qr_detected"), bool) - ): - raise DouyinError("Douyin login page origin is not allowed") - screenshot_params = {"format": "png", "fromSurface": False} - if isinstance(page.get("clip"), dict): - screenshot_params["clip"] = page["clip"] + deadline = time.monotonic() + LOGIN_RENDER_TIMEOUT + while True: + page = cdp.evaluate(qr_expression) + if ( + not isinstance(page, dict) + or page.get("origin") not in LOGIN_ORIGINS + or not isinstance(page.get("qr_detected"), bool) + ): + raise DouyinError("Douyin login page origin is not allowed") + if page["qr_detected"]: + if not isinstance(page.get("clip"), dict): + raise DouyinError("Douyin login QR code bounds are invalid") + break + if time.monotonic() >= deadline: + raise DouyinError("Douyin login QR code did not appear; check the browser login page") + time.sleep(0.2) + screenshot_params = {"format": "png", "fromSurface": True, "clip": page["clip"]} screenshot = cdp.command("Page.captureScreenshot", screenshot_params) if not isinstance(screenshot, dict): raise DouyinError("Douyin login screenshot is invalid") @@ -543,6 +559,11 @@ class DouyinBrowser: raise DouyinError("Douyin identity response is invalid") from exc user = payload.get("user") if isinstance(payload, dict) else None uid = str(user.get("uid", "")) if isinstance(user, dict) else "" + if response.status == 200 and isinstance(payload, dict) and ( + payload.get("status_code") == 8 + or (payload.get("status_code") == 0 and isinstance(user, dict) and uid == "0") + ): + raise DouyinLoginPending("Douyin login is awaiting user confirmation") if ( response.status != 200 or not isinstance(payload, dict) @@ -564,6 +585,8 @@ class DouyinBrowser: "unique_id": unique_id, "nickname": user.get("nickname", "") if isinstance(user, dict) else "", "short_id": user.get("short_id", "") if isinstance(user, dict) else "", + "avatar_url": next(iter((user.get("avatar_larger") or user.get("avatar_medium") or user.get("avatar_thumb") or {}).get("url_list") or []), ""), + "douyin_number": unique_id or str(user.get("short_id") or ""), } extra = payload.get("extra") if isinstance(payload, dict) else None server_now = extra.get("now") if isinstance(extra, dict) else None diff --git a/browser_gateway/runtime.py b/browser_gateway/runtime.py index 1c06664..250e918 100644 --- a/browser_gateway/runtime.py +++ b/browser_gateway/runtime.py @@ -4,6 +4,7 @@ from __future__ import annotations import fcntl import hashlib +import errno import http.client import json import logging @@ -201,15 +202,57 @@ class DisplayAllocator: ) -> _Lease: unavailable = unavailable or (lambda _value: False) for value in range(start, end + 1): - if unavailable(value) or (X11_SOCKET_DIR / f"X{value}").exists(): + if unavailable(value): continue lock = FileLock(self.lock_dir / f"display-{value}.lock") - if lock.acquire(): - if not (X11_SOCKET_DIR / f"X{value}").exists(): - return _Lease(value, lock) + if not lock.acquire(): + continue + try: + if self._in_use(value): + lock.release() + continue + if (X11_SOCKET_DIR / f"X{value}").exists(): + LOG.warning("reusing inactive X11 display :%d; Xvfb will replace stale files", value, extra={"display": value}) + return _Lease(value, lock) + except Exception: lock.release() + raise raise BrowserRuntimeError("no free Xvfb display is available", 503) + def _in_use(self, value: int) -> bool: + # Xvfb removes its own stale lock/socket at startup. File existence alone + # is not evidence of a live server after an unclean stop. + lock_path = X11_SOCKET_DIR.parent / f".X{value}-lock" + try: + pid = int(lock_path.read_text().strip()) + except FileNotFoundError: + pid = None + except ValueError as exc: + raise BrowserRuntimeError(f"invalid X11 PID lock for display {value}", 503) from exc + if pid is not None: + if pid <= 0: + raise BrowserRuntimeError(f"invalid X11 PID lock for display {value}", 503) + try: + os.kill(pid, 0) + return True + except ProcessLookupError: + pass # A dead server's files are reclaimed by Xvfb, not deleted here. + except PermissionError: + return True + path = str(X11_SOCKET_DIR / f"X{value}") + for address in (path, "\0" + path): + with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as probe: + probe.settimeout(0.2) + try: + probe.connect(address) + return True + except TimeoutError: + return True + except OSError as exc: + if exc.errno not in (errno.ENOENT, errno.ECONNREFUSED): + raise + return False + def reserve_existing(self, value: int) -> _Lease: if value < 1: raise BrowserRuntimeError("runtime display is invalid", 409) @@ -1104,6 +1147,9 @@ class NativeRuntimeManager: "--remote-allow-origins=http://127.0.0.1", "--disable-gpu", "--disable-gpu-compositing", + # Managed Xvfb sessions must not wait on a desktop keyring to + # initialize cookies; that blocks navigation before any request. + "--password-store=basic", f"--user-data-dir={record.profile_dir}", "--no-first-run", "--no-default-browser-check", @@ -1118,7 +1164,9 @@ class NativeRuntimeManager: flock = shutil.which("flock") if not flock: raise BrowserRuntimeError("flock is unavailable for runtime ownership", 503) - return [flock, str(lock_path), *command] + # Keep the actual server as systemd's MainPID so KillMode=mixed sends + # SIGTERM to Xvfb itself, allowing it to remove its display files. + return [flock, "--no-fork", str(lock_path), *command] def _profile_dir(self, profile_id: str) -> Path: if not isinstance(profile_id, str) or not PROFILE_ID_RE.fullmatch(profile_id): @@ -1418,6 +1466,7 @@ def validate_runtime_input(value: Mapping[str, Any]) -> None: "--remote-debugging-port", "--remote-allow-origins", "--remote-debugging-pipe", + "--password-store", "--proxy-server", "--display", "--headless", diff --git a/browser_gateway/server/http.py b/browser_gateway/server/http.py index f19750b..78a2609 100644 --- a/browser_gateway/server/http.py +++ b/browser_gateway/server/http.py @@ -22,7 +22,6 @@ from typing import cast from urllib.parse import parse_qs, urlsplit from ..platform.douyin import ( - ACCOUNT_KEY_RE, ACTIONS, COMMENTS_PATH, IDENTITY_URL, @@ -30,6 +29,7 @@ from ..platform.douyin import ( WORKS_PATH, DouyinBrowser, DouyinError, + DouyinLoginPending, SubscriptionManager, is_douyin_content_url, is_douyin_share_url, @@ -253,13 +253,16 @@ class Gateway: if ( not valid_douyin_generation(input) or not isinstance(expected_account_key, str) - or not ACCOUNT_KEY_RE.fullmatch(expected_account_key) + or (expected_account_key and not UID_RE.fullmatch(expected_account_key)) ): - raise RequestError("invalid Douyin identity request", 400) + raise RequestError("invalid Douyin identity request: expected account key must be a UID", 400) with self._alias_lock(alias): self._require_douyin_generation(alias, input) try: identity = self.browser.identity(alias) + except DouyinLoginPending: + LOG.info("Douyin identity awaiting_login alias=%s", alias) + return {"status": "manual_login", "reason": "awaiting_login"} except DouyinError as exc: LOG.warning( "Douyin identity verification failed alias=%s reason=%s", @@ -267,15 +270,11 @@ class Gateway: str(exc), ) raise RequestError( - "Douyin login identity could not be verified" + f"Douyin login identity could not be verified: {exc}" ) from exc - if expected_account_key not in { - identity["uid"], - identity["sec_uid"], - identity["unique_id"], - }: + if expected_account_key and expected_account_key != identity["uid"]: raise RequestError( - "Douyin identity does not match the expected account", 409 + f"已绑定 UID {expected_account_key},当前浏览器登录 UID {identity['uid']};请登录已绑定账号", 409 ) return identity diff --git a/browser_gateway/test_gateway.py b/browser_gateway/test_gateway.py index 2dd3650..71498f1 100644 --- a/browser_gateway/test_gateway.py +++ b/browser_gateway/test_gateway.py @@ -841,7 +841,7 @@ class BrowserTests(unittest.TestCase): cdp = Mock() cdp.evaluate.side_effect = [ True, - {"origin": "https://www.douyin.com", "qr_detected": True}, + {"origin": "https://www.douyin.com", "qr_detected": True, "clip": {"x": 0, "y": 0, "width": 180, "height": 180, "scale": 1}}, ] cdp.command.side_effect = [ {}, @@ -875,6 +875,35 @@ class BrowserTests(unittest.TestCase): self.assertLess(methods.index('wait_event'), methods.index('evaluate')) self.assertFalse(any("cookie" in expression.lower() for expression in cdp.evaluate.call_args.args)) + def test_login_qr_waits_for_actual_code_before_cropping(self) -> None: + screenshot = base64.b64encode(b'x' * 24000).decode('ascii') + cdp = Mock() + cdp.command.side_effect = [{}, {'frameId': 'frame-1'}, {'data': screenshot}, {'data': screenshot}] + clip = {'x': 10, 'y': 20, 'width': 180, 'height': 180, 'scale': 1} + cdp.evaluate.side_effect = [True, {'origin': 'https://www.douyin.com', 'qr_detected': False}, {'origin': 'https://www.douyin.com', 'qr_detected': True, 'clip': clip}] + browser = DouyinBrowser() + self._with_connection(browser, cast(BrowserCDP, cdp)) + with patch('browser_gateway.platform.douyin.time.sleep'): + result = browser.login_qr('safe') + self.assertTrue(result.qr_detected) + self.assertEqual(cdp.command.call_args.kwargs, {}) + self.assertEqual(cdp.command.call_args.args[1]['clip'], clip) + self.assertTrue(cdp.command.call_args.args[1]['fromSurface']) + self.assertEqual(cdp.evaluate.call_count, 3) + + def test_login_qr_reports_missing_code_instead_of_returning_page(self) -> None: + from .platform.douyin import LOGIN_RENDER_TIMEOUT + screenshot = base64.b64encode(b'x' * 24000).decode('ascii') + cdp = Mock() + cdp.command.side_effect = [{}, {'frameId': 'frame-1'}, {'data': screenshot}] + cdp.evaluate.side_effect = [False, {'origin': 'https://www.douyin.com', 'qr_detected': False}] + browser = DouyinBrowser() + self._with_connection(browser, cast(BrowserCDP, cdp)) + with patch('browser_gateway.platform.douyin.time.sleep'), patch('browser_gateway.platform.douyin.time.monotonic', side_effect=[0, 0, LOGIN_RENDER_TIMEOUT + 1]): + with self.assertRaisesRegex(DouyinError, 'login QR code did not appear'): + browser.login_qr('safe') + self.assertEqual([call.args[0] for call in cdp.command.call_args_list].count('Page.captureScreenshot'), 1) + def test_login_qr_does_not_capture_an_unloaded_page(self) -> None: cdp = Mock() cdp.command.side_effect = [{}, {"frameId": "frame-1"}] diff --git a/browser_gateway/test_identity_binding.py b/browser_gateway/test_identity_binding.py new file mode 100644 index 0000000..e1f4385 --- /dev/null +++ b/browser_gateway/test_identity_binding.py @@ -0,0 +1,34 @@ +"""First login discovers a UID; subsequent logins must match that UID.""" +import unittest +from contextlib import nullcontext +from unittest.mock import Mock, patch + +from .server.http import Gateway, RequestError + + +class IdentityBindingTests(unittest.TestCase): + def setUp(self): + self.browser = Mock() + self.browser.identity.return_value = {"uid": "123456789", "sec_uid": "MS4wLjAB-test", "unique_id": "douyin-name"} + self.gateway = Gateway(Mock(), "token", "node", browser=self.browser) + self.generation = {"binding_version": 1, "runtime_id": "a" * 64, "network_id": "native-" + "b" * 32} + self.lock = patch.object(self.gateway, "_alias_lock", return_value=nullcontext()) + self.check = patch.object(self.gateway, "_require_douyin_generation") + self.lock.start() + self.check_mock = self.check.start() + self.addCleanup(self.lock.stop) + self.addCleanup(self.check.stop) + + def test_first_login_reads_uid_without_a_preconfigured_key(self): + identity = self.gateway.douyin_identity("account-a", self.generation) + self.assertEqual(identity["uid"], "123456789") + self.check_mock.assert_called_once() + + def test_bound_login_requires_same_uid(self): + identity = self.gateway.douyin_identity("account-a", {**self.generation, "expected_account_key": "123456789"}) + self.assertEqual(identity["uid"], "123456789") + with self.assertRaises(RequestError) as caught: + self.gateway.douyin_identity("account-a", {**self.generation, "expected_account_key": "987654321"}) + self.assertEqual(caught.exception.status, 409) + self.assertIn("987654321", str(caught.exception)) + self.assertIn("123456789", str(caught.exception)) diff --git a/browser_gateway/test_login_pending.py b/browser_gateway/test_login_pending.py new file mode 100644 index 0000000..6eed92e --- /dev/null +++ b/browser_gateway/test_login_pending.py @@ -0,0 +1,57 @@ +"""Unit tests: an unsigned-in browser is pending, not a transport/identity error.""" +import json +import unittest +from contextlib import nullcontext +from types import SimpleNamespace +from unittest.mock import Mock + +from .platform.douyin import BrowserResponse, DouyinBrowser, DouyinError +from .server.http import Gateway, RequestError + + +class LoginPendingTests(unittest.TestCase): + def browser_with_response(self, body, status=200): + browser = DouyinBrowser() + browser.get = Mock(return_value=BrowserResponse(status, json.dumps(body), {})) + return browser + + def test_unsigned_browser_has_explicit_pending_signal(self): + from .platform.douyin import DouyinLoginPending + for body in [{"status_code": 8}, {"status_code": 0, "user": {"uid": "0"}}]: + with self.subTest(body=body): + browser = self.browser_with_response(body) + with self.assertRaises(DouyinLoginPending): + browser.identity("safe") + + def test_bad_response_is_not_pending(self): + from .platform.douyin import DouyinLoginPending + for body, status in [({}, 200), ({"status_code": 0}, 200), ({"status_code": 0, "user": {}}, 200), ({"status_code": 0, "user": {"uid": ""}}, 200), ({"status_code": 500}, 200), ({"status_code": 8}, 503), ({"status_code": 0, "user": {"uid": "bad"}}, 200)]: + with self.subTest(body=body, status=status): + with self.assertRaises(DouyinError) as raised: + self.browser_with_response(body, status).identity("safe") + self.assertNotIsInstance(raised.exception, DouyinLoginPending) + + def gateway(self): + runtimes = Mock() + runtimes.alias_lock.return_value = nullcontext() + runtimes.require_generation.return_value = SimpleNamespace(runtime_id='a' * 64, network_exit_id='exit-1') + browser = Mock() + gateway = Gateway(runtimes, 'gateway-token-123456', 'node-a', browser=browser) + payload = {'runtime_id': 'a' * 64, 'binding_version': 1, 'network_id': 'native-' + 'b' * 32, 'network_exit_id': 'exit-1', 'expected_account_key': '123'} + return gateway, browser, payload + + def test_gateway_returns_pending_as_normal_status_and_logs_it(self): + from .platform.douyin import DouyinLoginPending + gateway, browser, payload = self.gateway() + browser.identity.side_effect = DouyinLoginPending("waiting for login") + with self.assertLogs("creatorhub.gateway", level="INFO") as logs: + result = gateway.douyin_identity("safe", payload) + self.assertEqual(result, {"status": "manual_login", "reason": "awaiting_login"}) + self.assertIn("awaiting_login", " ".join(logs.output)) + + def test_gateway_keeps_real_failures_visible(self): + gateway, browser, payload = self.gateway() + browser.identity.side_effect = DouyinError("runtime disconnected") + with self.assertLogs("creatorhub.gateway", level="WARNING"): + with self.assertRaisesRegex(RequestError, "could not be verified"): + gateway.douyin_identity("safe", payload) diff --git a/browser_gateway/test_runtime.py b/browser_gateway/test_runtime.py index cdc4bd5..8b328a6 100644 --- a/browser_gateway/test_runtime.py +++ b/browser_gateway/test_runtime.py @@ -1,11 +1,15 @@ from __future__ import annotations import json +import errno +import os +import socket import subprocess import tempfile import unittest from pathlib import Path from typing import Any +from unittest.mock import patch from .runtime import ( BrowserRuntimeError, @@ -168,6 +172,67 @@ class SystemdUnitManagerTests(unittest.TestCase): displays.reserve_existing(0) +class DisplayAllocationTests(unittest.TestCase): + def test_dead_x_server_files_do_not_exhaust_display_pool(self) -> None: + from .runtime import DisplayAllocator + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + sockets = root / '.X11-unix' + sockets.mkdir() + stale = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + stale.bind(str(sockets / 'X100')) + stale.close() + (root / '.X100-lock').write_text('99999999\n') + with patch('browser_gateway.runtime.X11_SOCKET_DIR', sockets), patch('browser_gateway.runtime.os.kill', side_effect=ProcessLookupError(errno.ESRCH, 'gone')): + allocator = DisplayAllocator(root / 'leases') + lease = allocator.reserve(100, 100, lambda value: False) + self.assertEqual(lease.value, 100) + lease.lock.release() + + def test_invalid_pid_lock_is_reported_without_leaking_reservation(self) -> None: + from .runtime import DisplayAllocator + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + sockets = root / '.X11-unix' + sockets.mkdir() + pid_lock = root / '.X100-lock' + with patch('browser_gateway.runtime.X11_SOCKET_DIR', sockets): + allocator = DisplayAllocator(root / 'leases') + for value in ('bad', '0', '-1'): + with self.subTest(value=value): + pid_lock.write_text(value) + with self.assertRaisesRegex(BrowserRuntimeError, 'invalid X11 PID lock'): + allocator.reserve(100, 100, lambda value: False) + pid_lock.unlink() + lease = allocator.reserve(100, 100, lambda value: False) + self.assertEqual(lease.value, 100) + lease.lock.release() + + def test_live_x_server_pid_is_not_reused(self) -> None: + from .runtime import DisplayAllocator + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + sockets = root / '.X11-unix' + sockets.mkdir() + (root / '.X100-lock').write_text(str(os.getpid())) + with patch('browser_gateway.runtime.X11_SOCKET_DIR', sockets): + with self.assertRaisesRegex(BrowserRuntimeError, 'no free Xvfb display'): + DisplayAllocator(root / 'leases').reserve(100, 100, lambda value: False) + + def test_listening_x_socket_without_pid_lock_is_not_reused(self) -> None: + from .runtime import DisplayAllocator + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + sockets = root / '.X11-unix' + sockets.mkdir() + with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as server: + server.bind(str(sockets / 'X100')) + server.listen() + with patch('browser_gateway.runtime.X11_SOCKET_DIR', sockets): + with self.assertRaisesRegex(BrowserRuntimeError, 'no free Xvfb display'): + DisplayAllocator(root / 'leases').reserve(100, 100, lambda value: False) + + class NativeRuntimeManagerTests(unittest.TestCase): def setUp(self) -> None: self.temp = tempfile.TemporaryDirectory() @@ -188,6 +253,10 @@ class NativeRuntimeManagerTests(unittest.TestCase): self.manager._wait_for_display = lambda record: None self.manager._wait_for_cdp = lambda record: None + def test_runtime_lock_executes_server_as_main_process_for_graceful_stop(self) -> None: + command = self.manager._locked_command(self.state / 'test.lock', ['/bin/true']) + self.assertEqual(command[1], '--no-fork') + def tearDown(self) -> None: self.manager.close() self.temp.cleanup() @@ -220,6 +289,7 @@ class NativeRuntimeManagerTests(unittest.TestCase): self.assertNotIn("--", browser_command) self.assertIn("--disable-gpu", browser_command) self.assertIn("--disable-gpu-compositing", browser_command) + self.assertIn("--password-store=basic", browser_command) self.assertIn("--remote-debugging-address=127.0.0.1", browser_command) self.assertIn("--user-data-dir=" + stored["profile_dir"], browser_command) @@ -368,10 +438,9 @@ class NativeRuntimeManagerTests(unittest.TestCase): self.assertEqual(self.manager.list_public()[0]["state"], "released") def test_reserved_browser_flags_are_rejected_before_side_effects(self) -> None: - with self.assertRaises(BrowserRuntimeError): - self.manager.create( - self.payload(cmd=["--no-sandbox", "about:blank"]) - ) + for flag in ("--no-sandbox", "--password-store=gnome"): + with self.subTest(flag=flag), self.assertRaises(BrowserRuntimeError): + self.manager.create(self.payload(cmd=[flag, "about:blank"])) self.assertEqual(list(self.state.glob("runtimes/*/runtime.json")), []) def test_stopped_runtime_can_be_created_without_processes(self) -> None: diff --git a/internal/account/store.go b/internal/account/store.go index 947864a..9577279 100644 --- a/internal/account/store.go +++ b/internal/account/store.go @@ -154,7 +154,7 @@ func (s *Store) CreateAccount(ctx context.Context, account Account, credentials if _, err := tx.ExecContext(ctx, ` INSERT INTO social_account (account_id, credential_provider, credential_key, name, platform, platform_account_key, tags, status) - VALUES ($1, $2, $3, $4, $5, $6, $7, 'paused')`, account.ID, + VALUES ($1, $2, $3, $4, $5, NULLIF($6, ''), $7, 'paused')`, account.ID, account.CredentialReference.Provider, account.CredentialKey, account.Name, account.Platform, account.PlatformAccountKey, account.Tags); err != nil { return publicDatabaseError(err) @@ -187,7 +187,7 @@ func commitKnownRolledBack(err error) bool { func (s *Store) ListAccounts(ctx context.Context) ([]Account, error) { rows, err := s.db.QueryContext(ctx, ` - SELECT account.account_id, account.name, account.platform, account.platform_account_key, account.tags, + SELECT account.account_id, account.name, account.platform, COALESCE(account.platform_account_key, ''), account.tags, account.status, account.version FROM social_account account ORDER BY account.created_at, account.account_id`) @@ -211,7 +211,7 @@ func (s *Store) GetAccount(ctx context.Context, id string) (Account, error) { return Account{}, ErrInvalid } return scanAccount(s.db.QueryRowContext(ctx, ` - SELECT account.account_id, account.name, account.platform, account.platform_account_key, account.tags, + SELECT account.account_id, account.name, account.platform, COALESCE(account.platform_account_key, ''), account.tags, account.status, account.version FROM social_account account WHERE account.account_id = $1`, id)) @@ -255,7 +255,7 @@ func scanAccount(row accountScanner) (Account, error) { func validAccount(account Account) bool { if !idPattern.MatchString(account.ID) || strings.TrimSpace(account.Name) != account.Name || account.Name == "" || !utf8.ValidString(account.Name) || utf8.RuneCountInString(account.Name) > 128 || - !platformKeyPattern.MatchString(account.PlatformAccountKey) || len(account.Tags) > 20 || + (account.PlatformAccountKey != "" && !platformKeyPattern.MatchString(account.PlatformAccountKey)) || len(account.Tags) > 20 || len(account.Cookies) > 8192 || !refPattern.MatchString(account.CredentialReference.ID) || !credentialKeyPattern.MatchString(account.CredentialKey) || (account.CredentialReference.Provider != "os_keyring" && account.CredentialReference.Provider != "secret_manager") { diff --git a/internal/controlplane/api/account_environment_unit_test.go b/internal/controlplane/api/account_environment_unit_test.go index 5a88021..a8c5346 100644 --- a/internal/controlplane/api/account_environment_unit_test.go +++ b/internal/controlplane/api/account_environment_unit_test.go @@ -38,20 +38,8 @@ func TestAccountEnvironmentAutoBindingAndStart(t *testing.T) { app := fiber.New() RegisterAccountRoutes(app, accountStore, hubStore, &testCredentialBridge{values: map[string]string{}}) - // 网关缺失:创建账号即绑定失败 → 503 environment_binding_failed,但账号已存在(可补建重试) - response := do(app, http.MethodPost, "/api/phase-a/accounts", - `{"name":"测试账号","platform":"douyin","platform_account_key":"key-binding-1","cookies":"sessionid=1"}`) - if response.Code != http.StatusServiceUnavailable { - t.Fatalf("expected 503 when no gateway exists, got %d: %s", response.Code, response.Body.String()) - } - var failure map[string]string - if err := json.Unmarshal(response.Body.Bytes(), &failure); err != nil || failure["reason_code"] != "environment_binding_failed" || failure["account_id"] == "" { - t.Fatalf("binding failure payload: %s err=%v", response.Body.String(), err) - } - accountID := failure["account_id"] - if _, err := accountStore.GetAccount(ctx, accountID); err != nil { - t.Fatalf("account must exist after failed binding: %v", err) - } + // Existing-account maintenance remains available; public pre-login creation is removed. + accountID := maintenanceAccountFixture(t, ctx, accountStore) if response := do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/environment", ""); response.Code != http.StatusNotFound { t.Fatalf("expected 404 environment rebind without gateway, got %d: %s", response.Code, response.Body.String()) } @@ -66,7 +54,7 @@ func TestAccountEnvironmentAutoBindingAndStart(t *testing.T) { if _, err := hubStore.CreateGateway(ctx, "gw-main", gatewayServer.URL, gateway.token); err != nil { t.Fatal(err) } - response = do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/environment", "") + response := do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/environment", "") if response.Code != http.StatusOK { t.Fatalf("expected 200 environment rebind, got %d: %s", response.Code, response.Body.String()) } @@ -110,6 +98,17 @@ func TestAccountEnvironmentAutoBindingAndStart(t *testing.T) { t.Fatalf("start audit pair missing: rows=%d err=%v", startAudits, err) } + // 运行中的账号再次启动只核对已有运行实例,不创建新实例或报 runtime_active。 + runtimeID := environment.RuntimeID + response = do(app, http.MethodPost, "/api/phase-a/accounts/"+accountID+"/start", "") + if response.Code != http.StatusNoContent { + t.Fatalf("running account start must be idempotent: %d %s", response.Code, response.Body.String()) + } + environment, err = hubStore.GetEnvironmentContext(ctx, accountID) + if err != nil || environment.RuntimeID != runtimeID { + t.Fatalf("repeated start replaced the running runtime: %#v err=%v", environment, err) + } + } func TestAccountEnvironmentDeletionLifecycle(t *testing.T) { @@ -149,19 +148,10 @@ func TestAccountEnvironmentDeletionLifecycle(t *testing.T) { t.Fatalf("expected 404 deleting missing account, got %d: %s", response.Code, response.Body.String()) } - // 创建即绑定 → 审计与浏览器环境随账号生成 - response := do(app, http.MethodPost, "/api/phase-a/accounts", - `{"name":"待删账号","platform":"douyin","platform_account_key":"key-delete-1","cookies":"sessionid=1"}`) - if response.Code != http.StatusCreated { - t.Fatalf("expected 201 create account, got %d: %s", response.Code, response.Body.String()) + accountID := maintenanceAccountFixture(t, ctx, accountStore) + if _, _, err := ensureAccountEnvironment(ctx, hubStore, accountID, hub.Fingerprint{}); err != nil { + t.Fatal(err) } - var created struct { - ID string `json:"id"` - } - if err := json.Unmarshal(response.Body.Bytes(), &created); err != nil || created.ID == "" { - t.Fatalf("create payload: %s err=%v", response.Body.String(), err) - } - accountID := created.ID if _, err := hubStore.GetEnvironmentContext(ctx, accountID); err != nil { t.Fatalf("auto-bound environment missing: %v", err) } diff --git a/internal/controlplane/api/account_fingerprint_test.go b/internal/controlplane/api/account_fingerprint_test.go index 9eef9e5..1dbbc74 100644 --- a/internal/controlplane/api/account_fingerprint_test.go +++ b/internal/controlplane/api/account_fingerprint_test.go @@ -5,119 +5,94 @@ import ( "database/sql" "encoding/json" "net/http" - "net/http/httptest" "os" "testing" accountdomain "git.ipao.vip/rogee/creator-hub/internal/account" + "git.ipao.vip/rogee/creator-hub/internal/creator" hub "git.ipao.vip/rogee/creator-hub/internal/environment" "github.com/gofiber/fiber/v3" ) -// 创建社媒账号的指纹浏览器环境表单:自定义指纹随创建请求落库; -// seed 由服务端从账号派生、proxy 由出口体系管理(客户端值被忽略/清空)。 -func TestAccountCreateAcceptsFingerprintForm(t *testing.T) { +func maintenanceAccountFixture(t *testing.T, ctx context.Context, store *accountdomain.Store) string { + t.Helper() + id := accountdomain.NewAccountID() + account := accountdomain.Account{ID: id, Name: "已有账号", Platform: "douyin", PlatformAccountKey: id, Tags: []string{}, CredentialReference: accountdomain.CredentialReference{ID: id + "-cookies", Provider: "os_keyring"}, CredentialKey: "creatorhub/" + id + "/cookies", RuntimeStatus: "paused", Version: 1} + if err := store.CreateAccount(ctx, account, &testCredentialBridge{values: map[string]string{}}); err != nil { + t.Fatal(err) + } + return id +} + +func TestEnvironmentCreationDoesNotCreatePlatformAccount(t *testing.T) { databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") if databaseURL == "" { - t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage") + t.Skip("set CREATORHUB_POSTGRES_TEST_URL") } - ctx := context.Background() databaseURL = isolatedControlPlaneDatabaseURL(t, databaseURL) - accountStore, err := accountdomain.Open(ctx, databaseURL) + ctx := context.Background() + hs, err := hub.Open(ctx, databaseURL) if err != nil { t.Fatal(err) } - t.Cleanup(func() { _ = accountStore.Close() }) - hubStore, err := hub.Open(ctx, databaseURL) + t.Cleanup(func() { _ = hs.Close() }) + cs, err := creator.Open(ctx, databaseURL) if err != nil { t.Fatal(err) } - t.Cleanup(func() { _ = hubStore.Close() }) - gateway := &fakeGateway{token: "unit-test-gateway-token"} - gatewayServer := httptest.NewServer(gateway.handler(t)) - t.Cleanup(gatewayServer.Close) + t.Cleanup(func() { _ = cs.Close() }) + db, err := sql.Open("pgx", databaseURL) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = db.Close() }) app := fiber.New() - RegisterAccountRoutes(app, accountStore, hubStore, &testCredentialBridge{values: map[string]string{}}) - if _, err := hubStore.CreateGateway(ctx, "gw-main", gatewayServer.URL, gateway.token); err != nil { + registerEnvironmentLoginRoutes(app, cs, hs) + response := do(app, http.MethodPost, "/api/creator/environments", `{}`) + if response.Code != http.StatusNotFound { + t.Fatalf("missing gateway: %d %s", response.Code, response.Body.String()) + } + if _, err := hs.CreateGateway(ctx, "gw-main", "http://127.0.0.1:28187", "test-token-environment"); err != nil { t.Fatal(err) } - - response := do(app, http.MethodPost, "/api/phase-a/accounts", - `{"name":"指纹账号","platform":"douyin","platform_account_key":"key-fp-1","cookies":"sessionid=1", - "fingerprint":{"seed":42,"platform":"windows","platform_version":"10.0.0","brand":"Chrome","brand_version":"132.0.6834.159", - "hardware_concurrency":8,"lang":"zh-CN","accept_lang":"zh-CN,en-US","timezone":"Asia/Shanghai", - "proxy_server":"socks5://proxy.example:1080","disable_spoofing":"canvas,gpu"}}`) + response = do(app, http.MethodPost, "/api/creator/environments", `{"fingerprint":{"platform":"windows","lang":"zh-CN","timezone":"Asia/Shanghai"}}`) if response.Code != http.StatusCreated { - t.Fatalf("expected 201 create account with fingerprint, got %d: %s", response.Code, response.Body.String()) + t.Fatalf("create: %d %s", response.Code, response.Body.String()) } - var created struct { - ID string `json:"id"` + var env hub.EnvironmentContext + if err := json.Unmarshal(response.Body.Bytes(), &env); err != nil { + t.Fatal(err) } - if err := json.Unmarshal(response.Body.Bytes(), &created); err != nil || created.ID == "" { - t.Fatalf("create payload: %s err=%v", response.Body.String(), err) + if env.AccountID != "" || env.ProfileID != env.Alias || env.Fingerprint.Seed < 1001 || env.Fingerprint.Lang != "zh-CN" { + t.Fatalf("unbound environment: %#v", env) } - auditDB, err := sql.Open("pgx", databaseURL) + var count int + if err := db.QueryRowContext(ctx, `SELECT count(*) FROM social_account`).Scan(&count); err != nil || count != 0 { + t.Fatalf("creation made placeholder account: %d %v", count, err) + } + pending, err := hs.ListPendingEnvironments(ctx) + if err != nil || len(pending) != 1 { + t.Fatalf("pending: %v %v", pending, err) + } + for _, body := range []string{`{"name":"手填昵称"}`, `{"platform_account_key":"123"}`, `{"fingerprint":{"platform":"android"}}`, `{"fingerprint":{"seed":42}}`, `{"fingerprint":{"proxy_server":"http://invalid:8080"}}`} { + if response := do(app, http.MethodPost, "/api/creator/environments", body); response.Code != http.StatusBadRequest { + t.Fatalf("invalid input accepted: %s %d %s", body, response.Code, response.Body.String()) + } + } + result, err := cs.RecordVerifiedEnvironmentLogin(ctx, env.Alias, creator.PlatformIdentity{UID: "99491952055", Nickname: "同步昵称", DouyinNumber: "1004291301", AvatarURL: "https://example.com/avatar.jpg"}) if err != nil { t.Fatal(err) } - t.Cleanup(func() { _ = auditDB.Close() }) - var accountRowID int64 - if err := auditDB.QueryRowContext(ctx, `SELECT id FROM social_account WHERE account_id = $1`, created.ID).Scan(&accountRowID); err != nil { - t.Fatal(err) + bound, err := hs.GetEnvironmentContext(ctx, env.Alias) + if err != nil || bound.AccountID != result.AccountID || bound.ProfileID != env.ProfileID || bound.Fingerprint.Seed != env.Fingerprint.Seed { + t.Fatalf("binding replaced browser: %#v %v", bound, err) } - environment, err := hubStore.GetEnvironmentContext(ctx, created.ID) - if err != nil { - t.Fatal(err) + pending, err = hs.ListPendingEnvironments(ctx) + if err != nil || len(pending) != 0 { + t.Fatalf("bound env still pending: %v %v", pending, err) } - fingerprint := environment.Fingerprint - if fingerprint.Platform != "windows" || fingerprint.PlatformVersion != "10.0.0" || - fingerprint.Brand != "Chrome" || fingerprint.BrandVersion != "132.0.6834.159" || - fingerprint.HardwareConcurrency != 8 || fingerprint.Lang != "zh-CN" || - fingerprint.AcceptLang != "zh-CN,en-US" || fingerprint.Timezone != "Asia/Shanghai" || - fingerprint.DisableSpoofing != "canvas,gpu" { - t.Fatalf("fingerprint not persisted as submitted: %#v", fingerprint) - } - if fingerprint.Seed != accountRowID+1000 || fingerprint.ProxyServer != "" { - t.Fatalf("seed must be server-derived and proxy cleared: %#v (row id %d)", fingerprint, accountRowID) - } - - // 非法指纹值 → 400,账号不落库(校验前置,无创建后绑定失败的中间态)。 - if response := do(app, http.MethodPost, "/api/phase-a/accounts", - `{"name":"坏指纹","platform":"douyin","platform_account_key":"key-fp-2","fingerprint":{"platform":"android"}}`); response.Code != http.StatusBadRequest { - t.Fatalf("expected 400 for invalid fingerprint, got %d: %s", response.Code, response.Body.String()) - } - var invalidCount int - if err := auditDB.QueryRowContext(ctx, `SELECT count(*) FROM social_account WHERE platform_account_key = 'key-fp-2'`).Scan(&invalidCount); err != nil || invalidCount != 0 { - t.Fatalf("invalid fingerprint must not create account: rows=%d err=%v", invalidCount, err) - } - - // 幂等补建端点可携带同一指纹表单重试(创建时绑定失败的场景)。 - retry := do(app, http.MethodPost, "/api/phase-a/accounts", - `{"name":"补建账号","platform":"douyin","platform_account_key":"key-fp-3","fingerprint":{"timezone":"Asia/Shanghai"}}`) - if retry.Code != http.StatusCreated { - t.Fatalf("expected 201 create account for rebind retry, got %d: %s", retry.Code, retry.Body.String()) - } - var rebindCreated struct { - ID string `json:"id"` - } - if err := json.Unmarshal(retry.Body.Bytes(), &rebindCreated); err != nil || rebindCreated.ID == "" { - t.Fatalf("rebind create payload: %s err=%v", retry.Body.String(), err) - } - // 幂等补建端点可携带创建时未落库的指纹重试:先删除环境模拟"创建时绑定失败",补建后指纹落库。 - if err := hubStore.DeleteAccountEnvironment(ctx, rebindCreated.ID); err != nil { - t.Fatal(err) - } - rebind := do(app, http.MethodPost, "/api/phase-a/accounts/"+rebindCreated.ID+"/environment", - `{"fingerprint":{"platform":"linux","lang":"en-US"}}`) - if rebind.Code != http.StatusOK { - t.Fatalf("expected 200 environment rebind with fingerprint, got %d: %s", rebind.Code, rebind.Body.String()) - } - reboundEnvironment, err := hubStore.GetEnvironmentContext(ctx, rebindCreated.ID) - if err != nil { - t.Fatal(err) - } - // 补建以传入指纹为准(时区为空 = 创建时的时区不保留)。 - if reboundEnvironment.Fingerprint.Platform != "linux" || reboundEnvironment.Fingerprint.Lang != "en-US" || reboundEnvironment.Fingerprint.Timezone != "" { - t.Fatalf("rebind fingerprint mismatch: %#v", reboundEnvironment.Fingerprint) + RegisterAccountRoutes(app, nil, nil, nil) + if response := do(app, http.MethodPost, "/api/phase-a/accounts", `{"name":"废弃入口"}`); response.Code != http.StatusMethodNotAllowed { + t.Fatalf("old creation path remains: %d", response.Code) } } diff --git a/internal/controlplane/api/accounts_operations.go b/internal/controlplane/api/accounts_operations.go index 2698d14..0571054 100644 --- a/internal/controlplane/api/accounts_operations.go +++ b/internal/controlplane/api/accounts_operations.go @@ -6,7 +6,6 @@ import ( "errors" "io" "strconv" - "strings" "time" accountdomain "git.ipao.vip/rogee/creator-hub/internal/account" @@ -14,62 +13,8 @@ import ( "github.com/gofiber/fiber/v3" ) -type accountRequest struct { - Name string `json:"name"` - Platform string `json:"platform"` - PlatformAccountKey string `json:"platform_account_key"` - Tags []string `json:"tags"` - Cookies string `json:"cookies"` - Fingerprint hub.Fingerprint `json:"fingerprint"` -} - // RegisterAccountRoutes exposes account lifecycle routes (create with auto-binding/list/detail/补建/start/pause/resume/revoke + audit). func RegisterAccountRoutes(app *fiber.App, store *accountdomain.Store, runtimeStore HubStore, credentials accountdomain.CredentialBridge) { - app.Post("/api/phase-a/accounts", func(c fiber.Ctx) error { - var input accountRequest - if err := decodePhaseA(c, &input); err != nil { - return phaseAError(c, err) - } - // 指纹表单校验前置:非法值直接 400,避免账号已建、环境绑定失败的中间态。 - input.Fingerprint.ProxyServer = "" - input.Fingerprint.DisableNonProxiedUDP = false - if err := input.Fingerprint.Validate(); err != nil { - return c.Status(fiber.StatusBadRequest).JSON(map[string]string{"error": "fingerprint: " + err.Error()}) - } - tags := input.Tags - if tags == nil { - tags = []string{} - } - for index := range tags { - tags[index] = strings.TrimSpace(tags[index]) - } - accountID := accountdomain.NewAccountID() - account := accountdomain.Account{ - ID: accountID, Name: strings.TrimSpace(input.Name), Platform: strings.TrimSpace(input.Platform), - PlatformAccountKey: strings.TrimSpace(input.PlatformAccountKey), Tags: tags, Cookies: strings.TrimSpace(input.Cookies), - CredentialReference: accountdomain.CredentialReference{ID: accountID + "-cookies", Provider: "os_keyring"}, - CredentialKey: "creatorhub/" + accountID + "/cookies", - RuntimeStatus: "paused", Version: 1, - } - if err := store.CreateAccount(c.Context(), account, credentials); err != nil { - if errors.Is(err, accountdomain.ErrAccountCreationUnknown) { - return c.Status(fiber.StatusServiceUnavailable).JSON(map[string]string{ - "error": "account creation result is unknown", "reason_code": "account_creation_result_unknown", "account_id": accountID, - }) - } - return phaseAError(c, err) - } - // 账号即环境:创建即绑定(幂等)。绑定失败透传原因与账号 ID,客户端可用幂等补建端点重试。 - if runtimeStore != nil { - if _, _, err := ensureAccountEnvironment(c.Context(), runtimeStore, accountID, input.Fingerprint); err != nil { - return c.Status(fiber.StatusServiceUnavailable).JSON(map[string]string{ - "error": err.Error(), "reason_code": "environment_binding_failed", "account_id": accountID, - }) - } - } - return c.Status(fiber.StatusCreated).JSON(account) - }) - app.Get("/api/phase-a/accounts", func(c fiber.Ctx) error { accounts, err := store.ListAccounts(c.Context()) if err != nil { @@ -153,12 +98,19 @@ func RegisterAccountRoutes(app *fiber.App, store *accountdomain.Store, runtimeSt return hubError(c, err) } defer unlock() - if err := store.ResumeAccount(c.Context(), c.Params("id")); err != nil { - if errors.Is(err, accountdomain.ErrConflict) { - return accountResumeConflict(c, store, runtimeStore, c.Params("id")) - } + account, err := store.GetAccount(c.Context(), c.Params("id")) + if err != nil { return phaseAError(c, err) } + // 已启用的账号交由浏览器启动流程核对既有运行实例,无需再次恢复账号。 + if account.RuntimeStatus != "active" { + if err := store.ResumeAccount(c.Context(), c.Params("id")); err != nil { + if errors.Is(err, accountdomain.ErrConflict) { + return accountResumeConflict(c, store, runtimeStore, c.Params("id")) + } + return phaseAError(c, err) + } + } if err := startAccountEnvironment(c.Context(), runtimeStore, c.Params("id")); err != nil { return hubError(c, err) } diff --git a/internal/controlplane/api/app_migrated_test.go b/internal/controlplane/api/app_migrated_test.go index 76a389c..9a89adf 100644 --- a/internal/controlplane/api/app_migrated_test.go +++ b/internal/controlplane/api/app_migrated_test.go @@ -418,7 +418,7 @@ func TestPhaseAAccountRequestRejectsUnknownFields(t *testing.T) { t.Run(name, func(t *testing.T) { app := fiber.New() app.Post("/", func(c fiber.Ctx) error { - var input accountRequest + var input struct { Fingerprint hub.Fingerprint `json:"fingerprint"` } if err := decodePhaseA(c, &input); err != nil { return phaseAError(c, err) } diff --git a/internal/controlplane/api/creator.go b/internal/controlplane/api/creator.go index 8661092..909c7d1 100644 --- a/internal/controlplane/api/creator.go +++ b/internal/controlplane/api/creator.go @@ -59,6 +59,7 @@ func registerCreator(app *fiber.App, store *creator.Store, phaseAStore *accountd } func registerCreatorWithServices(app *fiber.App, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store, analyzer creator.ThemeAnalyzer) { + registerEnvironmentLoginRoutes(app, store, hubStore) // Platform records enter through the managed collector/listener, not a public // client-supplied write. The explicit test namespace is kept for isolated // contract tests and never participates in the production listener. @@ -667,7 +668,7 @@ func creatorError(c fiber.Ctx, err error) error { status, message = fiber.StatusConflict, creator.ErrUncertain.Error() } response := map[string]string{"error": message} - if errors.Is(err, creator.ErrUnavailable) && err.Error() != creator.ErrUnavailable.Error() { + if err.Error() != message { response["reason"] = err.Error() } return c.Status(status).JSON(response) @@ -697,23 +698,41 @@ func flattenActionEvidence(destination map[string]string, prefix string, value a } } +var errCreatorLoginPending = errors.New("browser is awaiting user login") + func (browser creatorGatewayBrowser) Identity(ctx context.Context, expectedKey string) (string, error) { + identity, err := browser.IdentityProfile(ctx, expectedKey) + return identity.UID, err +} + +func (browser creatorGatewayBrowser) IdentityProfile(ctx context.Context, expectedKey string) (creator.PlatformIdentity, error) { payload := gatewayGenerationPayload(browser.environment) payload["expected_account_key"] = expectedKey status, body, err := gatewayCall(ctx, browser.gateway, http.MethodPost, "/v1/browsers/"+url.PathEscape(browser.environment.Alias)+"/douyin/identity", payload, 30*time.Second) if err != nil { - return "", err + return creator.PlatformIdentity{}, err } if status != http.StatusOK { - return "", fmt.Errorf("douyin identity verification rejected with HTTP %d: %s", status, string(body)) + return creator.PlatformIdentity{}, fmt.Errorf("douyin identity verification rejected with HTTP %d: %s", status, string(body)) } var identity struct { - UID string `json:"uid"` + creator.PlatformIdentity + Status string `json:"status"` + Reason string `json:"reason"` } - if err := json.Unmarshal(body, &identity); err != nil || identity.UID == "" { - return "", errors.New("douyin identity response omitted uid") + if err := json.Unmarshal(body, &identity); err != nil { + return creator.PlatformIdentity{}, fmt.Errorf("decode douyin identity response: %w", err) } - return identity.UID, nil + if identity.Status != "" { + if identity.Status == "manual_login" && identity.Reason == "awaiting_login" && identity.UID == "" { + return creator.PlatformIdentity{}, errCreatorLoginPending + } + return creator.PlatformIdentity{}, errors.New("douyin identity response returned an unexpected login status") + } + if identity.UID == "" { + return creator.PlatformIdentity{}, errors.New("douyin identity response omitted uid") + } + return identity.PlatformIdentity, nil } func (browser creatorGatewayBrowser) MessageHistory(ctx context.Context, expectedUID, targetUID, cursor string, limit int) (douyinMessageHistory, error) { @@ -742,6 +761,9 @@ func startCreatorEnvironment(ctx context.Context, store HubStore, environment hu if store == nil { return creator.ErrUnavailable } + if !accountRunnable(environment) { + return fmt.Errorf("账号已暂停,请先启动账号,再显示登录二维码: %w", hub.ErrConflict) + } action := actionForEnvironment("start", environment) if err := store.AppendEnvironmentAction(ctx, "environment_action_requested", action); err != nil { return err @@ -812,7 +834,7 @@ func creatorLoginQRCode(ctx context.Context, store *creator.Store, phaseAStore * return nil, err } if account.Platform != creator.PlatformDouyin || profile.Platform != creator.PlatformDouyin || - profile.PlatformAccountKey == "" || account.PlatformAccountKey != profile.PlatformAccountKey { + account.PlatformAccountKey != profile.PlatformAccountKey { return nil, creator.ErrConflict } unlock, lockErr := lockAccountResources(ctx, hubStore, accountID) @@ -860,7 +882,7 @@ func verifyCreatorAccount(ctx context.Context, store *creator.Store, phaseAStore if err != nil { return creator.LoginResult{}, err } - if account.Platform != profile.Platform || account.Platform != creator.PlatformDouyin || profile.PlatformAccountKey == "" || account.PlatformAccountKey != profile.PlatformAccountKey { + if account.Platform != profile.Platform || account.Platform != creator.PlatformDouyin || account.PlatformAccountKey != profile.PlatformAccountKey { return creator.LoginResult{}, creator.ErrConflict } environment, err := hubStore.GetEnvironmentContextForAccount(ctx, accountID) @@ -874,11 +896,14 @@ func verifyCreatorAccount(ctx context.Context, store *creator.Store, phaseAStore if err != nil { return creator.LoginResult{}, fmt.Errorf("%w: gateway unavailable: %v", creator.ErrUnavailable, err) } - uid, identityErr := verifyCreatorPlatformIdentity(ctx, account.Platform, gateway, environment, profile.PlatformAccountKey) + identity, identityErr := (creatorGatewayBrowser{gateway: gateway, environment: environment}).IdentityProfile(ctx, profile.PlatformAccountKey) + if errors.Is(identityErr, errCreatorLoginPending) { + return creator.LoginResult{AccountID: accountID, Status: "manual_login", Reason: "awaiting_login", CheckedAt: time.Now().UTC()}, nil + } if identityErr != nil { return creator.LoginResult{}, fmt.Errorf("%w: verify the manually logged-in browser identity: %v", creator.ErrConflict, identityErr) } - result, err := store.RecordVerifiedLoginResult(ctx, accountID, uid) + result, err := store.RecordVerifiedEnvironmentLogin(ctx, environment.Alias, identity) if err != nil { return creator.LoginResult{}, fmt.Errorf("persist verified account identity: %w", err) } diff --git a/internal/controlplane/api/creator_login_pending_test.go b/internal/controlplane/api/creator_login_pending_test.go new file mode 100644 index 0000000..d3bd0b5 --- /dev/null +++ b/internal/controlplane/api/creator_login_pending_test.go @@ -0,0 +1,51 @@ +package api + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "testing" + + hub "git.ipao.vip/rogee/creator-hub/internal/environment" +) + +func TestCreatorGatewayIdentityPendingIsExplicit(t *testing.T) { + for _, tc := range []struct { + name, body string + pending bool + }{ + {"waiting", `{"status":"manual_login","reason":"awaiting_login"}`, true}, + {"malformed pending", `{"status":"manual_login","reason":"unknown"}`, false}, + {"pending with uid", `{"status":"manual_login","reason":"awaiting_login","uid":"123"}`, false}, + {"unknown status with uid", `{"status":"unknown","uid":"123"}`, false}, + {"empty identity", `{}`, false}, + {"malformed json", `{`, false}, + } { + t.Run(tc.name, func(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { _, _ = w.Write([]byte(tc.body)) })) + defer server.Close() + browser := creatorGatewayBrowser{gateway: hub.Gateway{Endpoint: server.URL}, environment: testRunnableEnvironment()} + _, err := browser.Identity(context.Background(), "123") + if err == nil { + t.Fatal("non-identity response was accepted") + } + if errors.Is(err, errCreatorLoginPending) != tc.pending { + t.Fatalf("pending=%v err=%v", tc.pending, err) + } + }) + } +} + +func TestCreatorGatewayIdentityRejectsHTTPFailureEvenWithPendingBody(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusServiceUnavailable) + _, _ = w.Write([]byte(`{"status":"manual_login","reason":"awaiting_login"}`)) + })) + defer server.Close() + browser := creatorGatewayBrowser{gateway: hub.Gateway{Endpoint: server.URL}, environment: testRunnableEnvironment()} + _, err := browser.Identity(context.Background(), "123") + if err == nil || errors.Is(err, errCreatorLoginPending) { + t.Fatalf("real HTTP failure treated as waiting: %v", err) + } +} diff --git a/internal/controlplane/api/creator_login_test.go b/internal/controlplane/api/creator_login_test.go index 778f64e..5016af2 100644 --- a/internal/controlplane/api/creator_login_test.go +++ b/internal/controlplane/api/creator_login_test.go @@ -1 +1,40 @@ package api + +import ( + "context" + "errors" + "strings" + "testing" + + "git.ipao.vip/rogee/creator-hub/internal/creator" + hub "git.ipao.vip/rogee/creator-hub/internal/environment" +) + +func TestLoginQRStartupReportsGatewayFailureAndKeepsActionRecord(t *testing.T) { + if err := startCreatorEnvironment(context.Background(), nil, hub.EnvironmentContext{}); !errors.Is(err, creator.ErrUnavailable) { + t.Fatalf("missing store must fail visibly: %v", err) + } + store := newMemoryStore() + gatewayErr := errors.New("gateway lookup failed") + store.gatewayFn = func(string) (hub.Gateway, error) { return hub.Gateway{}, gatewayErr } + environment := hub.EnvironmentContext{Env: hub.Env{Alias: "account-a", Gateway: "gw-a"}, AccountStatus: "active"} + err := startCreatorEnvironment(context.Background(), store, environment) + if !errors.Is(err, gatewayErr) { + t.Fatalf("startup hid gateway failure: %v", err) + } + if len(store.actions) != 2 || store.actions[1].Outcome != "failed" || store.actions[1].ReasonCode != "gateway_unavailable" { + t.Fatalf("startup failure must remain traceable: %#v", store.actions) + } +} + +func TestLoginQRRequiresExplicitAccountStart(t *testing.T) { + store := newMemoryStore() + environment := hub.EnvironmentContext{AccountStatus: "paused"} + err := startCreatorEnvironment(context.Background(), store, environment) + if !errors.Is(err, hub.ErrConflict) || !strings.Contains(err.Error(), "账号已暂停") { + t.Fatalf("paused login should explain how to proceed: %v", err) + } + if len(store.actions) != 0 { + t.Fatalf("paused login must not request browser start: %#v", store.actions) + } +} diff --git a/internal/controlplane/api/environment_login.go b/internal/controlplane/api/environment_login.go new file mode 100644 index 0000000..382c78f --- /dev/null +++ b/internal/controlplane/api/environment_login.go @@ -0,0 +1,122 @@ +package api + +import ( + "context" + "errors" + "fmt" + "time" + + "git.ipao.vip/rogee/creator-hub/internal/creator" + hub "git.ipao.vip/rogee/creator-hub/internal/environment" + "github.com/gofiber/fiber/v3" + "github.com/sirupsen/logrus" +) + +func registerEnvironmentLoginRoutes(app *fiber.App, store *creator.Store, hubStore *hub.Store) { + app.Get("/api/creator/environments", func(c fiber.Ctx) error { + items, err := hubStore.ListPendingEnvironments(c.Context()) + if err != nil { + return hubError(c, err) + } + return c.JSON(items) + }) + app.Post("/api/creator/environments", func(c fiber.Ctx) error { + var input struct { + Fingerprint hub.Fingerprint `json:"fingerprint"` + } + if err := decodePhaseA(c, &input); err != nil { + return phaseAError(c, err) + } + gateway, err := soleGateway(c.Context(), hubStore) + if err != nil { + return hubError(c, err) + } + env, err := hubStore.CreateStandaloneEnv(c.Context(), gateway.Name, input.Fingerprint) + if err != nil { + return hubError(c, err) + } + logrus.WithFields(logrus.Fields{"event_type": "login_environment_created", "alias": env.Alias, "gateway": gateway.Name}).Info("browser environment created without platform account") + return c.Status(fiber.StatusCreated).JSON(env) + }) + app.Post("/api/creator/environments/:alias/verify", func(c fiber.Ctx) error { + result, err := verifyLoginEnvironment(c.Context(), store, hubStore, c.Params("alias")) + if err != nil { + return creatorError(c, err) + } + return c.JSON(result) + }) + app.Post("/api/creator/environments/:alias/login-qr", func(c fiber.Ctx) error { + alias := c.Params("alias") + unlock, err := hubStore.LockResources(c.Context(), []string{alias}, nil) + if err != nil { + return hubError(c, err) + } + defer unlock() + env, err := hubStore.GetEnvironmentContext(c.Context(), alias) + if err != nil { + return hubError(c, err) + } + if err := startCreatorEnvironment(c.Context(), hubStore, env); err != nil { + return hubError(c, err) + } + env, err = hubStore.GetEnvironmentContext(c.Context(), alias) + if err != nil { + return hubError(c, err) + } + gateway, err := hubStore.GetGateway(c.Context(), env.Gateway) + if err != nil { + return hubError(c, err) + } + qr, err := (douyinGatewayBrowser{gateway: gateway, environment: env}).LoginQR(c.Context()) + if err != nil { + return creatorError(c, fmt.Errorf("%w: capture login QR: %v", creator.ErrUnavailable, err)) + } + return c.JSON(map[string]any{"status": "manual_login", "content_type": qr.ContentType, "image_base64": qr.BodyBase64, "qr_detected": qr.QRDetected, "expires_at": time.Now().UTC().Add(creatorLoginQRLifetime).Format(time.RFC3339)}) + }) +} + +func verifyLoginEnvironment(ctx context.Context, store *creator.Store, hubStore *hub.Store, alias string) (creator.LoginResult, error) { + unlock, err := hubStore.LockResources(ctx, []string{alias}, nil) + if err != nil { + return creator.LoginResult{}, err + } + defer unlock() + env, err := hubStore.GetEnvironmentContext(ctx, alias) + if err != nil { + return creator.LoginResult{}, err + } + if err := startCreatorEnvironment(ctx, hubStore, env); err != nil { + return creator.LoginResult{}, err + } + env, err = hubStore.GetEnvironmentContext(ctx, alias) + if err != nil { + return creator.LoginResult{}, err + } + expected := "" + if env.AccountID != "" { + p, err := store.GetAccountProfile(ctx, env.AccountID) + if err != nil { + return creator.LoginResult{}, err + } + expected = p.PlatformAccountKey + } + gateway, err := hubStore.GetGateway(ctx, env.Gateway) + if err != nil { + return creator.LoginResult{}, err + } + identity, err := (creatorGatewayBrowser{gateway: gateway, environment: env}).IdentityProfile(ctx, expected) + if errors.Is(err, errCreatorLoginPending) { + return creator.LoginResult{Status: "manual_login", Reason: "awaiting_login", CheckedAt: time.Now().UTC()}, nil + } + if err != nil { + return creator.LoginResult{}, fmt.Errorf("%w: browser identity: %v", creator.ErrConflict, err) + } + result, err := store.RecordVerifiedEnvironmentLogin(ctx, alias, identity) + entry := logrus.WithFields(logrus.Fields{"event_type": "environment_login_verified", "alias": alias, "uid": identity.UID, "account_id": result.AccountID}) + if err != nil { + entry.WithError(err).Warn("browser account binding rejected") + return creator.LoginResult{}, err + } + entry.Info("platform account bound and profile synchronized") + return result, nil +} diff --git a/internal/controlplane/api/environment_login_test.go b/internal/controlplane/api/environment_login_test.go new file mode 100644 index 0000000..f01452e --- /dev/null +++ b/internal/controlplane/api/environment_login_test.go @@ -0,0 +1,118 @@ +package api + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "strings" + "testing" + + "git.ipao.vip/rogee/creator-hub/internal/creator" + hub "git.ipao.vip/rogee/creator-hub/internal/environment" + "github.com/gofiber/fiber/v3" +) + +func TestPendingEnvironmentLoginRoutesKeepBrowserIdentity(t *testing.T) { + raw := os.Getenv("CREATORHUB_POSTGRES_TEST_URL") + if raw == "" { + t.Skip("set CREATORHUB_POSTGRES_TEST_URL") + } + dbURL := isolatedControlPlaneDatabaseURL(t, raw) + ctx := context.Background() + hs, err := hub.Open(ctx, dbURL) + if err != nil { + t.Fatal(err) + } + defer hs.Close() + cs, err := creator.Open(ctx, dbURL) + if err != nil { + t.Fatal(err) + } + defer cs.Close() + alias := "" + running := false + logged := false + creates := 0 + gateway := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + runtime := map[string]any{"id": "login-runtime", "alias": alias, "state": "running", "binding_version": 1, "network_id": "login-network", "node_id": "login-node", "proxy_ready": true} + switch { + case r.Method == http.MethodGet && r.URL.Path == "/v1/browsers": + items := []map[string]any{} + if running { + items = append(items, runtime) + } + json.NewEncoder(w).Encode(items) + case r.Method == http.MethodGet && r.URL.Path == "/v1/browsers/"+alias: + if !running { + w.WriteHeader(http.StatusNotFound) + return + } + json.NewEncoder(w).Encode(runtime) + case r.Method == http.MethodPost && r.URL.Path == "/v1/browsers": + var body map[string]any + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + t.Error(err) + } + if body["profile_id"] != alias || body["alias"] != alias { + t.Errorf("unstable browser identity: %#v", body) + } + running = true + creates++ + w.WriteHeader(http.StatusCreated) + json.NewEncoder(w).Encode(runtime) + case strings.HasSuffix(r.URL.Path, "/douyin/identity"): + if !logged { + json.NewEncoder(w).Encode(map[string]string{"status": "manual_login", "reason": "awaiting_login"}) + return + } + json.NewEncoder(w).Encode(creator.PlatformIdentity{UID: "99491952055", Nickname: "平台昵称", AvatarURL: "https://example.com/avatar.jpg", DouyinNumber: "1004291301", SecUID: "MS4w-test"}) + case strings.HasSuffix(r.URL.Path, "/douyin/login-qr"): + json.NewEncoder(w).Encode(map[string]any{"content_type": "image/png", "body_base64": "aGk=", "qr_detected": true}) + default: + t.Errorf("unexpected native request %s %s", r.Method, r.URL.Path) + http.NotFound(w, r) + } + })) + defer gateway.Close() + if _, err := hs.CreateGateway(ctx, "login-gw", gateway.URL, "test-token-login"); err != nil { + t.Fatal(err) + } + env, err := hs.CreateStandaloneEnv(ctx, "login-gw", hub.Fingerprint{}) + if err != nil { + t.Fatal(err) + } + alias = env.Alias + app := fiber.New() + registerEnvironmentLoginRoutes(app, cs, hs) + route := "/api/creator/environments/" + alias + for _, suffix := range []string{"verify", "login-qr"} { + response := do(app, http.MethodPost, route+"/"+suffix, `{}`) + if response.Code != http.StatusOK { + t.Fatalf("%s: %d %s", suffix, response.Code, response.Body.String()) + } + } + pending, err := hs.ListPendingEnvironments(ctx) + if err != nil || len(pending) != 1 { + t.Fatalf("pending browser created account: %v %v", pending, err) + } + logged = true + response := do(app, http.MethodPost, route+"/verify", `{}`) + if response.Code != http.StatusOK { + t.Fatalf("verify: %d %s", response.Code, response.Body.String()) + } + var result creator.LoginResult + if err := json.Unmarshal(response.Body.Bytes(), &result); err != nil || result.AccountID == "" || result.Status != "logged_in" { + t.Fatalf("login response: %#v %v", result, err) + } + bound, err := hs.GetEnvironmentContext(ctx, alias) + if err != nil || bound.AccountID != result.AccountID || bound.RuntimeID != "login-runtime" || bound.ProfileID != alias || bound.Fingerprint.Seed != env.Fingerprint.Seed { + t.Fatalf("binding changed browser: %#v %v", bound, err) + } + response = do(app, http.MethodPost, route+"/verify", `{}`) + if response.Code != http.StatusOK || creates != 1 { + t.Fatalf("repeat login recreated browser: creates=%d %d %s", creates, response.Code, response.Body.String()) + } +} diff --git a/internal/controlplane/api/environments.go b/internal/controlplane/api/environments.go index 7ab9c06..6bdf95e 100644 --- a/internal/controlplane/api/environments.go +++ b/internal/controlplane/api/environments.go @@ -101,10 +101,7 @@ func gatewayCreatePayload(environment hub.EnvironmentContext, _ string, networkE fingerprint := environment.Fingerprint fingerprint.ProxyServer = "" fingerprint.DisableNonProxiedUDP = false - profileID := environment.AccountID - if profileID == "" { - profileID = environment.Alias - } + profileID := environment.ProfileID payload := map[string]any{ "alias": environment.Alias, "name": environment.Name, @@ -239,7 +236,7 @@ func gatewayRejected(status int, body []byte) error { if status >= 400 && status < 500 { return gatewayFailure{status: status, message: "gateway rejected: " + errorFromBody(body, status)} } - return gatewayFailure{status: http.StatusBadGateway, message: fmt.Sprintf("gateway call failed with status %d", status)} + return gatewayFailure{status: http.StatusBadGateway, message: fmt.Sprintf("gateway call failed with status %d: %s", status, errorFromBody(body, status))} } func gatewayUnreachable(err error) error { @@ -391,7 +388,7 @@ func validCreatedRuntime(created runtimeStatus, environment hub.EnvironmentConte } func accountRunnable(environment hub.EnvironmentContext) bool { - return environment.AccountStatus == "active" + return (environment.AccountID == "" && environment.AccountStatus == "") || environment.AccountStatus == "active" } func releaseRuntime(ctx context.Context, store RuntimeStopStore, environment hub.EnvironmentContext) error { @@ -1556,7 +1553,11 @@ func startBrowserRuntime(ctx context.Context, store HubStore, probe NetworkExitP reconcileErr := reconcileGatewayCreate(ctx, store, gateway, environment, status, callErr, body) if reconcileErr != nil { _ = finish("unknown", "gateway_result_unknown", environment) - return errors.Join(gatewayFailure{status: http.StatusBadGateway, message: "gateway start result unknown; retry to reconcile"}, reconcileErr) + cause := callErr + if cause == nil { + cause = gatewayRejected(status, body) + } + return errors.Join(gatewayFailure{status: http.StatusBadGateway, message: "gateway start result unknown; retry to reconcile"}, cause, reconcileErr) } _ = finish("failed", "gateway_create_failed", environment) if callErr != nil { diff --git a/internal/controlplane/api/gateway_error_reason_test.go b/internal/controlplane/api/gateway_error_reason_test.go new file mode 100644 index 0000000..91d3c7e --- /dev/null +++ b/internal/controlplane/api/gateway_error_reason_test.go @@ -0,0 +1,118 @@ +package api + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/http/httptest" + "strings" + "testing" + + hub "git.ipao.vip/rogee/creator-hub/internal/environment" +) + +func TestStartUnknownPreservesGatewayFailureReason(t *testing.T) { + for _, disconnected := range []bool{false, true} { + t.Run(map[bool]string{false: "server failure", true: "connection lost"}[disconnected], func(t *testing.T) { + created := false + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodPost { + created = true + if disconnected { + connection, _, err := w.(http.Hijacker).Hijack() + if err != nil { + t.Error(err) + return + } + _ = connection.Close() + return + } + w.WriteHeader(http.StatusServiceUnavailable) + _, _ = w.Write([]byte(`{"error":"no free Xvfb display is available"}`)) + return + } + if created { + w.WriteHeader(http.StatusServiceUnavailable) + _, _ = w.Write([]byte(`{"error":"cannot reconcile runtime"}`)) + return + } + _, _ = w.Write([]byte(`[]`)) + })) + defer server.Close() + store := newMemoryStore() + store.gateways["gw-1"] = hub.Gateway{Name: "gw-1", Endpoint: server.URL} + environment := hub.EnvironmentContext{Env: hub.Env{Alias: "account-a", Gateway: "gw-1", Fingerprint: hub.Fingerprint{Seed: 1}}, AccountID: "account-a", AccountStatus: "active", BindingVersion: 1} + store.bindings[environment.Alias] = environment + finish := func(outcome, reason string, _ hub.EnvironmentContext) error { + if outcome != "unknown" || reason != "gateway_result_unknown" { + t.Errorf("unexpected result: %s/%s", outcome, reason) + } + return nil + } + err := startBrowserRuntime(context.Background(), store, nil, environment, finish) + want := "no free Xvfb display is available" + if disconnected { + want = "EOF" + } + if err == nil || !strings.Contains(err.Error(), want) { + t.Fatalf("original startup cause was hidden: %v", err) + } + }) + } +} + +func TestStartBrowserRuntimeSuccessfulAndFailedCompletion(t *testing.T) { + for _, existing := range []bool{false, true} { + for _, finishFails := range []bool{false, true} { + t.Run(map[bool]string{false: "create", true: "reconcile"}[existing]+map[bool]string{false: " success", true: " finish error"}[finishFails], func(t *testing.T) { + posts := 0 + runtime := runtimeStatus{ID: "runtime-a", Alias: "account-a", State: "running", NetworkID: "net-a", BindingVersion: 1, ProxyReady: true} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodPost { + posts++ + w.WriteHeader(http.StatusCreated) + _ = json.NewEncoder(w).Encode(runtime) + } else if existing { + _ = json.NewEncoder(w).Encode([]runtimeStatus{runtime}) + } else { + _, _ = w.Write([]byte(`[]`)) + } + })) + defer server.Close() + store := newMemoryStore() + store.gateways["gw-1"] = hub.Gateway{Name: "gw-1", Endpoint: server.URL} + environment := hub.EnvironmentContext{Env: hub.Env{Alias: "account-a", Gateway: "gw-1", Fingerprint: hub.Fingerprint{Seed: 1}}, AccountID: "account-a", AccountStatus: "active", BindingVersion: 1} + store.bindings[environment.Alias] = environment + finishErr := errors.New("completion recording failed") + finish := func(outcome, reason string, current hub.EnvironmentContext) error { + want := "environment_started" + if existing { + want = "gateway_reconciled" + } + if outcome != "succeeded" || reason != want || current.RuntimeID != runtime.ID { + t.Errorf("unexpected completion: %s/%s/%s", outcome, reason, current.RuntimeID) + } + if finishFails { + return finishErr + } + return nil + } + err := startBrowserRuntime(context.Background(), store, nil, environment, finish) + if finishFails && !errors.Is(err, finishErr) || !finishFails && err != nil { + t.Fatalf("unexpected start result: %v", err) + } + if existing && posts != 0 || !existing && posts != 1 { + t.Fatalf("unexpected create count: %d", posts) + } + }) + } + } +} + +func TestGatewayRejectedPreservesServerFailureReason(t *testing.T) { + err := gatewayRejected(http.StatusServiceUnavailable, []byte(`{"error":"no free Xvfb display is available"}`)) + if !strings.Contains(err.Error(), "no free Xvfb display is available") { + t.Fatalf("gateway failure hid its cause: %v", err) + } +} diff --git a/internal/creator/account_uid_binding_test.go b/internal/creator/account_uid_binding_test.go new file mode 100644 index 0000000..bafb8ee --- /dev/null +++ b/internal/creator/account_uid_binding_test.go @@ -0,0 +1,121 @@ +package creator + +import ( + "errors" + "fmt" + "os" + "strings" + "sync" + "testing" + + "git.ipao.vip/rogee/creator-hub/internal/account" +) + +func TestCreatorFirstLoginBindsUIDAndNeverReplacesIt(t *testing.T) { + store, accounts, ctx := openCreatorIntegrationStore(t) + create := func(id string) { + t.Helper() + err := accounts.CreateAccount(ctx, account.Account{ID: id, Name: id, Platform: "douyin", CredentialReference: account.CredentialReference{ID: id, Provider: "secret_manager"}, CredentialKey: "creatorhub/" + id}, integrationCredentialBridge{}) + if err != nil { + t.Fatal(err) + } + if err := store.EnsureAccountProfile(ctx, id); err != nil { + t.Fatal(err) + } + } + create("uid-first") + create("uid-second") + before, err := store.GetAccountProfile(ctx, "uid-first") + if err != nil || before.PlatformAccountKey != "" { + t.Fatalf("new account must be unbound: %#v err=%v", before, err) + } + for i := 0; i < 2; i++ { + result, err := store.RecordVerifiedLoginResult(ctx, "uid-first", "99491952055") + if err != nil || result.Status != "logged_in" || result.ActualKey != "99491952055" { + t.Fatalf("verify %d: %#v err=%v", i, result, err) + } + } + bound, err := accounts.GetAccount(ctx, "uid-first") + if err != nil || bound.PlatformAccountKey != "99491952055" { + t.Fatalf("verified UID was not persisted: %#v err=%v", bound, err) + } + if _, err := store.RecordVerifiedLoginResult(ctx, "uid-first", "123456789"); !errors.Is(err, ErrConflict) || !strings.Contains(err.Error(), "99491952055") || !strings.Contains(err.Error(), "123456789") { + t.Fatalf("mismatch must explain both UIDs: %v", err) + } + if _, err := store.RecordVerifiedLoginResult(ctx, "uid-second", "99491952055"); !errors.Is(err, ErrConflict) || !strings.Contains(err.Error(), "已绑定") { + t.Fatalf("duplicate binding should fail visibly: %v", err) + } + second, err := store.GetAccountProfile(ctx, "uid-second") + if err != nil || second.PlatformAccountKey != "" || second.LoginStatus == "logged_in" { + t.Fatalf("failed bind changed account: %#v err=%v", second, err) + } +} + +func TestCreatorConcurrentFirstLoginBindsExactlyOneUID(t *testing.T) { + store, accounts, ctx := openCreatorIntegrationStore(t) + id := "uid-concurrent" + if err := accounts.CreateAccount(ctx, account.Account{ID: id, Name: id, Platform: "douyin", CredentialReference: account.CredentialReference{ID: id, Provider: "secret_manager"}, CredentialKey: "creatorhub/" + id}, integrationCredentialBridge{}); err != nil { + t.Fatal(err) + } + if err := store.EnsureAccountProfile(ctx, id); err != nil { + t.Fatal(err) + } + var wg sync.WaitGroup + outcomes := make(chan error, 2) + for _, uid := range []string{"111111111", "222222222"} { + wg.Add(1) + go func(uid string) { + defer wg.Done() + _, err := store.RecordVerifiedLoginResult(ctx, id, uid) + outcomes <- err + }(uid) + } + wg.Wait() + close(outcomes) + succeeded, conflicts := 0, 0 + for err := range outcomes { + if err == nil { + succeeded++ + } else if errors.Is(err, ErrConflict) { + conflicts++ + } else { + t.Fatal(err) + } + } + if succeeded != 1 || conflicts != 1 { + t.Fatalf("competing logins changed UID more than once: success=%d conflict=%d", succeeded, conflicts) + } +} + +func TestUIDMigrationClearsUnverifiedKeysAndPreservesVerifiedUID(t *testing.T) { + store, accounts, ctx := openCreatorIntegrationStore(t) + for i, status := range []string{"logged_in", "needs_login", "manual_required"} { + id := fmt.Sprintf("uid-migration-%d", i) + if err := accounts.CreateAccount(ctx, account.Account{ID: id, Name: id, Platform: "douyin", PlatformAccountKey: fmt.Sprintf("1000000%d", i), CredentialReference: account.CredentialReference{ID: id, Provider: "secret_manager"}, CredentialKey: "creatorhub/" + id}, integrationCredentialBridge{}); err != nil { + t.Fatal(err) + } + if err := store.EnsureAccountProfile(ctx, id); err != nil { + t.Fatal(err) + } + if _, err := store.db.ExecContext(ctx, `UPDATE creator_account_profile SET login_status=$2 WHERE account_id=(SELECT id FROM social_account WHERE account_id=$1)`, id, status); err != nil { + t.Fatal(err) + } + } + migration, err := os.ReadFile("../environment/migrations/1047_account_uid_first_login.sql") + if err != nil { + t.Fatal(err) + } + if _, err := store.db.ExecContext(ctx, string(migration)); err != nil { + t.Fatal(err) + } + for i := 0; i < 3; i++ { + profile, err := store.GetAccountProfile(ctx, fmt.Sprintf("uid-migration-%d", i)) + want := "" + if i == 0 { + want = "10000000" + } + if err != nil || profile.PlatformAccountKey != want { + t.Fatalf("migration %d: %#v err=%v", i, profile, err) + } + } +} diff --git a/internal/creator/accounts.go b/internal/creator/accounts.go index bd7c5d3..c65c611 100644 --- a/internal/creator/accounts.go +++ b/internal/creator/accounts.go @@ -23,12 +23,12 @@ func (s *Store) EnsureAccountProfile(ctx context.Context, accountID string) erro func accountProfileQuery() string { return ` - SELECT a.account_id, a.name, a.platform, a.platform_account_key, + SELECT a.account_id, a.name, a.platform, COALESCE(a.platform_account_key, ''), a.status, p.login_username, p.password_configured, p.real_name_status, p.real_name, p.identity_number, p.note, p.business_status, p.big_account, p.reply_requirements, p.login_status, p.login_reason, - p.login_checked_at, p.cooldown_seconds, p.updated_at + p.login_checked_at, p.cooldown_seconds, p.updated_at, p.avatar_url, p.douyin_number FROM social_account a JOIN creator_account_profile p ON p.account_id = a.id WHERE a.account_id = $1` @@ -43,7 +43,7 @@ func scanAccountProfile(scanner interface{ Scan(...any) error }) (AccountProfile &result.LoginUsername, &result.PasswordConfigured, &result.RealNameStatus, &result.RealName, &result.IdentityNumber, &result.Note, &result.BusinessStatus, &result.BigAccount, &result.ReplyRequirements, &result.LoginStatus, &result.LoginReason, - &checkedAt, &result.CooldownSeconds, &result.UpdatedAt, + &checkedAt, &result.CooldownSeconds, &result.UpdatedAt, &result.AvatarURL, &result.DouyinNumber, ); err != nil { return AccountProfile{}, err } @@ -404,30 +404,57 @@ func (s *Store) RecordLoginResult(ctx context.Context, accountID, status, reason func (s *Store) RecordVerifiedLoginResult(ctx context.Context, accountID, actualKey string) (LoginResult, error) { actualKey = strings.TrimSpace(actualKey) - if actualKey == "" { + if actualKey == "" || len(actualKey) > 128 { return LoginResult{}, ErrInvalid } - return s.recordLoginResult(ctx, accountID, "logged_in", "", actualKey) + if err := s.EnsureAccountProfile(ctx, accountID); err != nil { + return LoginResult{}, err + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return LoginResult{}, fmt.Errorf("begin account UID verification: %w", err) + } + defer tx.Rollback() + var id int64 + var boundUID sql.NullString + if err := tx.QueryRowContext(ctx, `SELECT id, platform_account_key FROM social_account WHERE account_id = $1 FOR UPDATE`, accountID).Scan(&id, &boundUID); err != nil { + return LoginResult{}, databaseError(err) + } + if boundUID.Valid && boundUID.String != actualKey { + return LoginResult{}, fmt.Errorf("%w: 已绑定 UID %s,当前浏览器登录 UID %s;请登录已绑定账号", ErrConflict, boundUID.String, actualKey) + } + if !boundUID.Valid { + if _, err := tx.ExecContext(ctx, `UPDATE social_account SET platform_account_key = $2, version = version + 1, updated_at = now() WHERE id = $1`, id, actualKey); err != nil { + if errors.Is(databaseError(err), ErrConflict) { + return LoginResult{}, fmt.Errorf("%w: UID %s 已绑定到其他账号,请使用已有账号记录", ErrConflict, actualKey) + } + return LoginResult{}, databaseError(err) + } + } + now := time.Now().UTC() + if _, err := tx.ExecContext(ctx, `UPDATE creator_account_profile SET login_status = 'logged_in', login_reason = '', login_checked_at = $2, updated_at = $2 WHERE account_id = $1`, id, now); err != nil { + return LoginResult{}, databaseError(err) + } + if err := tx.Commit(); err != nil { + return LoginResult{}, fmt.Errorf("commit account UID verification: %w", err) + } + return LoginResult{AccountID: accountID, Status: "logged_in", ActualKey: actualKey, CheckedAt: now}, nil } func (s *Store) recordLoginResult(ctx context.Context, accountID, status, reason, actualKey string) (LoginResult, error) { status = strings.TrimSpace(status) reason = strings.TrimSpace(reason) - if status != "logged_in" && status != "needs_login" && status != "failed" && status != "manual_required" { + if status != "needs_login" && status != "failed" && status != "manual_required" { return LoginResult{}, ErrInvalid } if len(actualKey) > 128 || utf8.RuneCountInString(reason) > 1000 { return LoginResult{}, ErrInvalid } - profile, err := s.GetAccountProfile(ctx, accountID) - if err != nil { + if _, err := s.GetAccountProfile(ctx, accountID); err != nil { return LoginResult{}, err } - if status == "logged_in" && (actualKey == "" || actualKey != profile.PlatformAccountKey) { - return LoginResult{}, ErrConflict - } now := time.Now().UTC() - _, err = s.db.ExecContext(ctx, ` + _, err := s.db.ExecContext(ctx, ` UPDATE creator_account_profile SET login_status = $2, login_reason = $3, login_checked_at = $4, updated_at = $4 WHERE account_id = (SELECT account.id FROM social_account account WHERE account.account_id = $1)`, accountID, status, reason, now) diff --git a/internal/creator/environment_login.go b/internal/creator/environment_login.go new file mode 100644 index 0000000..853f259 --- /dev/null +++ b/internal/creator/environment_login.go @@ -0,0 +1,80 @@ +package creator + +import ( + "context" + "database/sql" + "errors" + "fmt" + "strings" + "time" + "unicode/utf8" +) + +// PlatformIdentity contains only platform-supplied identity/profile fields. +// Notes, credentials and business settings are never copied from this response. +type PlatformIdentity struct { + UID string `json:"uid"` + Nickname string `json:"nickname"` + SecUID string `json:"sec_uid"` + AvatarURL string `json:"avatar_url"` + DouyinNumber string `json:"douyin_number"` +} + +// RecordVerifiedEnvironmentLogin binds a verified browser to its account in one +// transaction. The browser's profile, fingerprint and runtime are not replaced. +func (s *Store) RecordVerifiedEnvironmentLogin(ctx context.Context, alias string, identity PlatformIdentity) (LoginResult, error) { + identity.UID = strings.TrimSpace(identity.UID) + identity.Nickname = strings.TrimSpace(identity.Nickname) + if alias == "" || identity.UID == "" || identity.Nickname == "" || utf8.RuneCountInString(identity.Nickname) > 64 { + return LoginResult{}, fmt.Errorf("%w: 登录结果缺少有效 UID 或昵称", ErrInvalid) + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return LoginResult{}, fmt.Errorf("begin verified environment login: %w", err) + } + defer tx.Rollback() + var rowID sql.NullInt64 + if err := tx.QueryRowContext(ctx, `SELECT account_id FROM browser_env WHERE alias=$1 FOR UPDATE`, alias).Scan(&rowID); err != nil { + return LoginResult{}, rowError(err) + } + now := time.Now().UTC() + var accountID string + if rowID.Valid { + var expected string + if err := tx.QueryRowContext(ctx, `SELECT account_id,COALESCE(platform_account_key,'') FROM social_account WHERE id=$1 FOR UPDATE`, rowID.Int64).Scan(&accountID, &expected); err != nil { + return LoginResult{}, rowError(err) + } + if expected != "" && expected != identity.UID { + return LoginResult{}, fmt.Errorf("%w: 已绑定抖音 UID %s,当前浏览器登录的是 UID %s;不能覆盖已有绑定", ErrConflict, expected, identity.UID) + } + if _, err := tx.ExecContext(ctx, `UPDATE social_account SET platform_account_key=$2,name=$3,version=version+CASE WHEN platform_account_key IS DISTINCT FROM $2 THEN 1 ELSE 0 END,updated_at=$4 WHERE id=$1`, rowID.Int64, identity.UID, identity.Nickname, now); err != nil { + return LoginResult{}, identityWriteError(err, identity.UID) + } + } else { + accountID = newID("account") + if err := tx.QueryRowContext(ctx, `INSERT INTO social_account(account_id,name,platform,platform_account_key,credential_provider,credential_key,status) VALUES($1,$2,'douyin',$3,'os_keyring',$4,'active') RETURNING id`, accountID, identity.Nickname, identity.UID, "creatorhub/"+accountID+"/cookies").Scan(&rowID); err != nil { + return LoginResult{}, identityWriteError(err, identity.UID) + } + if _, err := tx.ExecContext(ctx, `UPDATE browser_env SET account_id=$2,updated_at=$3 WHERE alias=$1`, alias, rowID.Int64, now); err != nil { + return LoginResult{}, databaseError(err) + } + } + if _, err := tx.ExecContext(ctx, `INSERT INTO creator_account_profile(account_id) VALUES($1) ON CONFLICT(account_id) DO NOTHING`, rowID.Int64); err != nil { + return LoginResult{}, databaseError(err) + } + if _, err := tx.ExecContext(ctx, `UPDATE creator_account_profile SET login_status='logged_in',login_reason='',login_checked_at=$2,updated_at=$2,sec_uid=$3,avatar_url=$4,douyin_number=$5 WHERE account_id=$1`, rowID.Int64, now, identity.SecUID, identity.AvatarURL, identity.DouyinNumber); err != nil { + return LoginResult{}, databaseError(err) + } + if err := tx.Commit(); err != nil { + return LoginResult{}, identityWriteError(err, identity.UID) + } + return LoginResult{AccountID: accountID, Status: "logged_in", ActualKey: identity.UID, CheckedAt: now}, nil +} + +func identityWriteError(err error, uid string) error { + mapped := databaseError(err) + if errors.Is(mapped, ErrConflict) { + return fmt.Errorf("%w: 抖音 UID %s 已绑定其他账号;不能重复创建或覆盖已有绑定", ErrConflict, uid) + } + return mapped +} diff --git a/internal/creator/environment_login_test.go b/internal/creator/environment_login_test.go new file mode 100644 index 0000000..820cecb --- /dev/null +++ b/internal/creator/environment_login_test.go @@ -0,0 +1,111 @@ +package creator + +import ( + "errors" + "strings" + "sync" + "testing" +) + +func TestEnvironmentLoginCreatesAccountOnlyAfterVerification(t *testing.T) { + s, _, ctx := openCreatorIntegrationStore(t) + var gatewayID int64 + if err := s.db.QueryRowContext(ctx, `INSERT INTO gateway(name,endpoint,token) VALUES('browser-login-gw','http://127.0.0.1:28187','test-token-browser-login') RETURNING id`).Scan(&gatewayID); err != nil { + t.Fatal(err) + } + for i, alias := range []string{"env-login-a", "env-login-b"} { + if _, err := s.db.ExecContext(ctx, `INSERT INTO browser_env(alias,name,gateway_id,fingerprint,profile_id) VALUES($1,$1,$2,jsonb_build_object('seed',$3::bigint),$1)`, alias, gatewayID, 1001+i); err != nil { + t.Fatal(err) + } + } + var before int + if err := s.db.QueryRowContext(ctx, `SELECT count(*) FROM social_account`).Scan(&before); err != nil { + t.Fatal(err) + } + identity := PlatformIdentity{UID: "99491952055", Nickname: "平台昵称", SecUID: "MS4w-test", DouyinNumber: "1004291301", AvatarURL: "https://example.com/avatar.jpg"} + result, err := s.RecordVerifiedEnvironmentLogin(ctx, "env-login-a", identity) + if err != nil || result.Status != "logged_in" || result.AccountID == "" { + t.Fatalf("first login: %#v %v", result, err) + } + profile, err := s.GetAccountProfile(ctx, result.AccountID) + if err != nil || profile.Name != identity.Nickname || profile.PlatformAccountKey != identity.UID || profile.AvatarURL != identity.AvatarURL || profile.DouyinNumber != identity.DouyinNumber { + t.Fatalf("synced profile: %#v %v", profile, err) + } + var after int + var profileID string + var seed int64 + if err := s.db.QueryRowContext(ctx, `SELECT count(*) FROM social_account`).Scan(&after); err != nil { + t.Fatal(err) + } + if after != before+1 { + t.Fatalf("account count: before=%d after=%d", before, after) + } + if err := s.db.QueryRowContext(ctx, `SELECT profile_id,(fingerprint->>'seed')::bigint FROM browser_env WHERE alias='env-login-a'`).Scan(&profileID, &seed); err != nil { + t.Fatal(err) + } + if profileID != "env-login-a" || seed != 1001 { + t.Fatal("binding changed browser identity") + } + if _, err := s.db.ExecContext(ctx, `UPDATE creator_account_profile SET note='本地备注',big_account=true WHERE account_id=(SELECT id FROM social_account WHERE account_id=$1)`, result.AccountID); err != nil { + t.Fatal(err) + } + identity.Nickname = "改名后的昵称" + again, err := s.RecordVerifiedEnvironmentLogin(ctx, "env-login-a", identity) + if err != nil || again.AccountID != result.AccountID { + t.Fatalf("repeat login: %#v %v", again, err) + } + profile, err = s.GetAccountProfile(ctx, result.AccountID) + if err != nil || profile.Name != identity.Nickname || profile.Note != "本地备注" || !profile.BigAccount { + t.Fatalf("metadata refresh changed local settings: %#v %v", profile, err) + } + wrong := identity + wrong.UID = "123456789" + if _, err := s.RecordVerifiedEnvironmentLogin(ctx, "env-login-a", wrong); !errors.Is(err, ErrConflict) || !strings.Contains(err.Error(), identity.UID) || !strings.Contains(err.Error(), wrong.UID) { + t.Fatalf("mismatch: %v", err) + } + if _, err := s.RecordVerifiedEnvironmentLogin(ctx, "env-login-b", identity); !errors.Is(err, ErrConflict) || !strings.Contains(err.Error(), "已绑定") { + t.Fatalf("duplicate: %v", err) + } + var bound bool + if err := s.db.QueryRowContext(ctx, `SELECT account_id IS NOT NULL FROM browser_env WHERE alias='env-login-b'`).Scan(&bound); err != nil || bound { + t.Fatalf("duplicate created binding: %v %v", bound, err) + } +} + +func TestConcurrentEnvironmentLoginsCreateExactlyOneAccount(t *testing.T) { + s, _, ctx := openCreatorIntegrationStore(t) + var gatewayID int64 + if err := s.db.QueryRowContext(ctx, `INSERT INTO gateway(name,endpoint,token) VALUES('concurrent-login-gw','http://127.0.0.1:28187','test-token-concurrent') RETURNING id`).Scan(&gatewayID); err != nil { + t.Fatal(err) + } + for i, alias := range []string{"concurrent-env-a", "concurrent-env-b"} { + if _, err := s.db.ExecContext(ctx, `INSERT INTO browser_env(alias,name,gateway_id,fingerprint,profile_id) VALUES($1,$1,$2,jsonb_build_object('seed',$3::bigint),$1)`, alias, gatewayID, 1101+i); err != nil { + t.Fatal(err) + } + } + var wg sync.WaitGroup + outcomes := make(chan error, 2) + for _, alias := range []string{"concurrent-env-a", "concurrent-env-b"} { + wg.Add(1) + go func(alias string) { + defer wg.Done() + _, err := s.RecordVerifiedEnvironmentLogin(ctx, alias, PlatformIdentity{UID: "777777777", Nickname: "并发账号"}) + outcomes <- err + }(alias) + } + wg.Wait() + close(outcomes) + success, conflict := 0, 0 + for err := range outcomes { + if err == nil { + success++ + } else if errors.Is(err, ErrConflict) { + conflict++ + } else { + t.Fatal(err) + } + } + if success != 1 || conflict != 1 { + t.Fatalf("success=%d conflict=%d", success, conflict) + } +} diff --git a/internal/creator/models.go b/internal/creator/models.go index f322579..438b170 100644 --- a/internal/creator/models.go +++ b/internal/creator/models.go @@ -56,6 +56,8 @@ type AccountProfile struct { Name string `json:"name"` Platform string `json:"platform"` PlatformAccountKey string `json:"platform_account_key"` + AvatarURL string `json:"avatar_url"` + DouyinNumber string `json:"douyin_number"` RuntimeStatus string `json:"runtime_status"` LoginUsername string `json:"login_username"` PasswordConfigured bool `json:"password_configured"` diff --git a/internal/environment/environment.go b/internal/environment/environment.go index 5846522..ec836dc 100644 --- a/internal/environment/environment.go +++ b/internal/environment/environment.go @@ -47,6 +47,7 @@ type ExitObservation struct { type EnvironmentContext struct { Env + ProfileID string `json:"profile_id"` AccountID string `json:"account_id"` AccountStatus string `json:"account_status"` BindingVersion int64 `json:"binding_version"` @@ -425,8 +426,8 @@ func (s *Store) CreateBoundEnv(ctx context.Context, env Env, accountID, exitID s } var created string if err := tx.QueryRowContext(ctx, ` - INSERT INTO browser_env (alias, name, gateway_id, fingerprint, account_id, exit_id, version) - SELECT $1, $2, gateway.id, jsonb_set($4::jsonb, '{seed}', to_jsonb(account.id + 1000)), account.id, network.id, 1 + INSERT INTO browser_env (alias, name, gateway_id, fingerprint, account_id, exit_id, version, profile_id) + SELECT $1, $2, gateway.id, jsonb_set($4::jsonb, '{seed}', to_jsonb(nextval('browser_fingerprint_seed_seq'))), account.id, network.id, 1, $1 FROM social_account account JOIN gateway ON gateway.name = $3 LEFT JOIN network_exit network ON network.exit_id = NULLIF($6, '') @@ -462,7 +463,7 @@ func (s *Store) GetEnvironmentContext(ctx context.Context, alias string) (Enviro var cleanupBindingVersion sql.NullInt64 err = tx.QueryRowContext(ctx, ` SELECT environment.alias, environment.name, gateway.name, - environment.fingerprint, environment.created_at, account.account_id, account.status, + environment.fingerprint, environment.created_at, COALESCE(account.account_id, ''), COALESCE(account.status, ''), environment.profile_id, environment.version, environment.runtime_cleanup_pending, environment.runtime_cleanup_binding_version, environment.runtime_cleanup_runtime_id, environment.runtime_cleanup_network_id, @@ -474,12 +475,12 @@ func (s *Store) GetEnvironmentContext(ctx context.Context, alias string) (Enviro COALESCE(network.created_at, to_timestamp(0)), COALESCE(network.updated_at, to_timestamp(0)), COALESCE(environment.runtime_id, ''), COALESCE(environment.runtime_network_id, ''), COALESCE(environment.runtime_node_id, '') FROM browser_env environment - JOIN social_account account ON account.id = environment.account_id + LEFT JOIN social_account account ON account.id = environment.account_id JOIN gateway ON gateway.id = environment.gateway_id LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1`, alias). Scan(&result.Alias, &result.Name, &result.Gateway, &encoded, &result.CreatedAt, - &result.AccountID, &result.AccountStatus, + &result.AccountID, &result.AccountStatus, &result.ProfileID, &result.BindingVersion, &result.RuntimeCleanupPending, &cleanupBindingVersion, &cleanupRuntimeID, &cleanupNetworkID, &result.Exit.ID, &result.Exit.Protocol, &result.Exit.Host, &result.Exit.Port, @@ -730,17 +731,22 @@ func (s *Store) ActivateRuntime(ctx context.Context, alias, runtimeID string, bi var currentBindingVersion int64 var cleanupPending bool err = tx.QueryRowContext(ctx, ` - SELECT account.account_id, environment.version, COALESCE(network.exit_id, ''), environment.runtime_cleanup_pending, - account.status + SELECT COALESCE(account.account_id, ''), environment.version, COALESCE(network.exit_id, ''), environment.runtime_cleanup_pending, + COALESCE(account.status, '') FROM browser_env environment - JOIN social_account account ON account.id = environment.account_id + LEFT JOIN social_account account ON account.id = environment.account_id LEFT JOIN network_exit network ON network.id = environment.exit_id - WHERE environment.alias = $1 FOR UPDATE OF environment, account`, alias). + WHERE environment.alias = $1 FOR UPDATE OF environment`, alias). Scan(&accountID, ¤tBindingVersion, ¤tExitID, &cleanupPending, &accountStatus) if err != nil { return EnvironmentContext{}, rowError(err) } - if cleanupPending || accountStatus != "active" || + if accountID != "" { + if err := tx.QueryRowContext(ctx, `SELECT status FROM social_account WHERE account_id=$1 FOR UPDATE`, accountID).Scan(&accountStatus); err != nil { + return EnvironmentContext{}, publicDatabaseError(err) + } + } + if cleanupPending || (accountID != "" && accountStatus != "active") || currentBindingVersion != bindingVersion || currentExitID != exitID { return EnvironmentContext{}, ErrConflict } @@ -791,12 +797,12 @@ func (s *Store) ReleaseRuntime(ctx context.Context, environment EnvironmentConte var accountID, exitID string var bindingVersion int64 err = tx.QueryRowContext(ctx, ` - SELECT account.account_id, COALESCE(network.exit_id, ''), environment.version + SELECT COALESCE(account.account_id, ''), COALESCE(network.exit_id, ''), environment.version FROM browser_env environment - JOIN social_account account ON account.id = environment.account_id + LEFT JOIN social_account account ON account.id = environment.account_id LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1 AND environment.version = $2 AND environment.runtime_id = $3 - FOR UPDATE OF environment, account`, + FOR UPDATE OF environment`, environment.Alias, environment.BindingVersion, environment.RuntimeID). Scan(&accountID, &exitID, &bindingVersion) if errors.Is(err, sql.ErrNoRows) { @@ -837,11 +843,11 @@ func (s *Store) SetRuntimeCleanupPending(ctx context.Context, environment Enviro if err := tx.QueryRowContext(ctx, ` SELECT environment.runtime_cleanup_pending, environment.runtime_cleanup_binding_version, environment.runtime_cleanup_runtime_id, environment.runtime_cleanup_network_id, - account.account_id, COALESCE(network.exit_id, '') + COALESCE(account.account_id, ''), COALESCE(network.exit_id, '') FROM browser_env environment - JOIN social_account account ON account.id = environment.account_id + LEFT JOIN social_account account ON account.id = environment.account_id LEFT JOIN network_exit network ON network.id = environment.exit_id - WHERE environment.alias = $1 AND environment.version = $2 FOR UPDATE OF environment, account`, environment.Alias, environment.BindingVersion). + WHERE environment.alias = $1 AND environment.version = $2 FOR UPDATE OF environment`, environment.Alias, environment.BindingVersion). Scan(¤tPending, &cleanupBindingVersion, &cleanupRuntimeID, &cleanupNetworkID, &accountID, &exitID); errors.Is(err, sql.ErrNoRows) { return ErrConflict diff --git a/internal/environment/migration043_probe_test.go b/internal/environment/migration043_probe_test.go index febb324..34eeb96 100644 --- a/internal/environment/migration043_probe_test.go +++ b/internal/environment/migration043_probe_test.go @@ -70,6 +70,8 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) { // audit_event 任务列已删 assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() AND table_name = 'audit_event' AND column_name IN ('confirmation_id','confirmation_version','attempt_id','task_id','runtime_instance_id')`, 0) - // 统一登记表:1-38(除 36)、1017-1045、43 全部登记 - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53) + // 统一登记表包含后续的首次登录 UID 绑定迁移。 + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 55) + assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema() + AND table_name = 'social_account' AND column_name = 'platform_account_key' AND is_nullable = 'YES'`, 1) } diff --git a/internal/environment/migration_test.go b/internal/environment/migration_test.go index 19e079e..c932302 100644 --- a/internal/environment/migration_test.go +++ b/internal/environment/migration_test.go @@ -30,7 +30,7 @@ func TestUnifiedAccountMigration(t *testing.T) { t.Fatal(err) } defer db.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 55) 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')`, 3) assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name = 'browser_image'`, 0) @@ -42,7 +42,7 @@ func TestUnifiedAccountMigration(t *testing.T) { store = openFullyMigratedHub(t, ctx, testURL) store.Close() - assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 53) + assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 55) }) t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) { diff --git a/internal/environment/migrations/1047_account_uid_first_login.sql b/internal/environment/migrations/1047_account_uid_first_login.sql new file mode 100644 index 0000000..e1feb8a --- /dev/null +++ b/internal/environment/migrations/1047_account_uid_first_login.sql @@ -0,0 +1,12 @@ +-- UID is unknown until the first verified browser login. NULL allows multiple +-- unbound accounts while retaining the existing (platform, UID) uniqueness. +ALTER TABLE social_account ALTER COLUMN platform_account_key DROP NOT NULL; + +-- Hand-entered identifiers are not proof of ownership. Preserve only bindings +-- that have already passed browser identity verification. +UPDATE social_account AS account +SET platform_account_key = NULL +WHERE NOT EXISTS ( + SELECT 1 FROM creator_account_profile AS profile + WHERE profile.account_id = account.id AND profile.login_status = 'logged_in' +); diff --git a/internal/environment/migrations/1048_environment_first_login.sql b/internal/environment/migrations/1048_environment_first_login.sql new file mode 100644 index 0000000..fbbc4f0 --- /dev/null +++ b/internal/environment/migrations/1048_environment_first_login.sql @@ -0,0 +1,29 @@ +-- Browser identity exists before a platform account. Binding must not change +-- either the fingerprint or the on-disk profile of an existing browser. +ALTER TABLE browser_env ALTER COLUMN account_id DROP NOT NULL; +ALTER TABLE browser_env ADD COLUMN profile_id text; +UPDATE browser_env e SET profile_id = COALESCE(a.account_id, e.alias) +FROM social_account a WHERE a.id = e.account_id; +UPDATE browser_env SET profile_id = alias WHERE profile_id IS NULL; +ALTER TABLE browser_env ALTER COLUMN profile_id SET NOT NULL; +-- A profile identifier is assigned at environment creation, not at login. +CREATE FUNCTION assign_browser_profile_id() RETURNS trigger LANGUAGE plpgsql AS $$ +BEGIN + IF TG_OP = 'INSERT' AND NEW.profile_id IS NULL THEN + NEW.profile_id := NEW.alias; + END IF; + IF TG_OP = 'UPDATE' AND NEW.profile_id IS DISTINCT FROM OLD.profile_id THEN + RAISE EXCEPTION 'browser profile identity cannot be changed'; + END IF; + RETURN NEW; +END; +$$; +CREATE TRIGGER browser_profile_identity BEFORE INSERT OR UPDATE ON browser_env +FOR EACH ROW EXECUTE FUNCTION assign_browser_profile_id(); + +CREATE SEQUENCE browser_fingerprint_seed_seq AS bigint MINVALUE 1001 START 1001; +SELECT setval('browser_fingerprint_seed_seq', GREATEST(1001, COALESCE(MAX((fingerprint->>'seed')::bigint),1000)), COALESCE(MAX((fingerprint->>'seed')::bigint),1000)>=1001) FROM browser_env; + +ALTER TABLE creator_account_profile ADD COLUMN sec_uid text NOT NULL DEFAULT ''; +ALTER TABLE creator_account_profile ADD COLUMN avatar_url text NOT NULL DEFAULT ''; +ALTER TABLE creator_account_profile ADD COLUMN douyin_number text NOT NULL DEFAULT ''; diff --git a/internal/environment/standalone.go b/internal/environment/standalone.go new file mode 100644 index 0000000..8e80dad --- /dev/null +++ b/internal/environment/standalone.go @@ -0,0 +1,54 @@ +package environment + +import ( + "context" + "encoding/json" + "fmt" +) + +// CreateStandaloneEnv creates only a browser environment; no social account or +// login profile exists until the browser's first verified login. +func (s *Store) CreateStandaloneEnv(ctx context.Context, gateway string, fingerprint Fingerprint) (EnvironmentContext, error) { + if !gatewayNamePattern.MatchString(gateway) || fingerprint.Seed != 0 || fingerprint.ProxyServer != "" || fingerprint.Validate() != nil { + return EnvironmentContext{}, ErrInvalid + } + alias := "env-" + newHubID() + encoded, err := json.Marshal(fingerprint) + if err != nil { + return EnvironmentContext{}, fmt.Errorf("encode browser fingerprint: %w", err) + } + if err := s.db.QueryRowContext(ctx, `INSERT INTO browser_env(alias,name,gateway_id,fingerprint,profile_id) SELECT $1,$1,id,jsonb_set($3::jsonb,'{seed}',to_jsonb(nextval('browser_fingerprint_seed_seq'))),$1 FROM gateway WHERE name=$2 RETURNING alias`, alias, gateway, encoded).Scan(&alias); err != nil { + return EnvironmentContext{}, rowError(err) + } + return s.GetEnvironmentContext(ctx, alias) +} + +func (s *Store) ListPendingEnvironments(ctx context.Context) ([]EnvironmentContext, error) { + rows, err := s.db.QueryContext(ctx, `SELECT alias FROM browser_env WHERE account_id IS NULL ORDER BY created_at,id`) + if err != nil { + return nil, fmt.Errorf("list pending browser environments: %w", err) + } + aliases := []string{} + for rows.Next() { + var alias string + if err := rows.Scan(&alias); err != nil { + rows.Close() + return nil, err + } + aliases = append(aliases, alias) + } + err = rows.Err() + rows.Close() + if err != nil { + return nil, err + } + result := make([]EnvironmentContext, 0, len(aliases)) + for _, alias := range aliases { + env, err := s.GetEnvironmentContext(ctx, alias) + if err != nil { + return nil, err + } + result = append(result, env) + } + return result, nil +} diff --git a/internal/environment/store.go b/internal/environment/store.go index b537ad4..4d5e4ee 100644 --- a/internal/environment/store.go +++ b/internal/environment/store.go @@ -179,6 +179,12 @@ var migration1045 string //go:embed migrations/1046_competitor_profile_signature.sql var migration1046 string +//go:embed migrations/1047_account_uid_first_login.sql +var migration1047 string + +//go:embed migrations/1048_environment_first_login.sql +var migration1048 string + var ( ErrConflict = errors.New("resource conflicts with existing state") ErrInvalid = errors.New("invalid environment input") @@ -319,7 +325,7 @@ func (s *Store) migrate(ctx context.Context) error { {1029, migration1029}, {1030, migration1030}, {1031, migration1031}, {1032, migration1032}, {1033, migration1033}, {1034, migration1034}, {1035, migration1035}, {1036, migration1036}, {1037, migration1037}, {1038, migration1038}, {1039, migration1039}, {1040, migration1040}, {1041, migration1041}, {1042, migration1042}, {1043, migration1043}, {1044, migration1044}, {1045, migration1045}, {1046, migration1046}, - {43, migration043}, {44, migration044}} { + {43, migration043}, {44, migration044}, {1047, migration1047}, {1048, migration1048}} { 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 environment schema migration state") diff --git a/internal/environment/store_test.go b/internal/environment/store_test.go index 3b6fb78..07e1030 100644 --- a/internal/environment/store_test.go +++ b/internal/environment/store_test.go @@ -106,10 +106,11 @@ func TestResourceLocksReserveConnectionsForLifecycleQueries(t *testing.T) { if databaseURL == "" { t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage") } + // 迁移初始化不属于连接锁的时限;只限时测量实际资源锁与生命周期查询。 + store := openFullyMigratedHub(t, context.Background(), isolatedDatabaseURL(t, databaseURL)) + t.Cleanup(func() { _ = store.Close() }) ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() - store := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL)) - t.Cleanup(func() { _ = store.Close() }) const workers = 10 var acquired atomic.Int32 diff --git a/web/src/components/accounts/AccountLoginModal.tsx b/web/src/components/accounts/AccountLoginModal.tsx new file mode 100644 index 0000000..908a3b4 --- /dev/null +++ b/web/src/components/accounts/AccountLoginModal.tsx @@ -0,0 +1,74 @@ +import { Alert, Button, Flex, Image, Modal, Spin, Typography } from 'antd'; +import { ReloadOutlined } from '@ant-design/icons'; +import React, { useEffect, useRef, useState } from 'react'; +import { request } from '@/requestErrorConfig'; +import { startAccountLogin, type LoginState } from './account-login'; + +const phaseTitles: Record = { + loading: '正在获取二维码', checking: '正在核验登录状态', waiting: '等待扫码或手机确认', + refreshing: '二维码已过期,正在自动刷新', success: '登录成功', error: '登录流程已停止', +}; + +export default function AccountLoginModal({ accountID, accountName, resource = 'accounts', onClose, onSuccess }: { + accountID: string; + accountName: string; + resource?: 'accounts' | 'environments'; + onClose: () => void; + onSuccess: () => void; +}) { + const [state, setState] = useState({ phase: 'loading', logs: [] }); + const [attempt, setAttempt] = useState(0); + const [now, setNow] = useState(Date.now()); + const successRef = useRef(onSuccess); + successRef.current = onSuccess; + useEffect(() => startAccountLogin(accountID, request, setState, () => successRef.current(), resource), [accountID, resource, attempt]); + useEffect(() => { + if (state.phase === 'success' || state.phase === 'error') return; + const timer = setInterval(() => setNow(Date.now()), 1000); + return () => clearInterval(timer); + }, [state.phase]); + const remaining = state.qr ? Math.max(0, Math.ceil((Date.parse(state.qr.expires_at) - now) / 1000)) : 0; + const busy = ['loading', 'checking', 'refreshing'].includes(state.phase); + const elapsed = Math.max(0, Math.floor((now - (state.logs.at(-1)?.time ?? now)) / 1000)); + return ( + } onClick={() => setAttempt(value => value + 1)}>重试登录 : null, + , + ]} + > + + + {state.phase !== 'success' && state.phase !== 'error' && ( + + {state.qr && remaining > 0 ? ( + 抖音登录二维码 + ) : busy ? ( + + + + {remaining === 0 && state.qr ? '二维码已到期,正在检查并刷新' : phaseTitles[state.phase]} · 已等待 {elapsed} 秒 + + + ) : null} + {state.qr && remaining > 0 ? 剩余 {remaining} 秒 · 到期自动刷新 : null} + {state.phase === 'checking' && remaining > 0 ? 正在检查登录状态 · 已等待 {elapsed} 秒 : null} + + )} + 登录进度 +
+ + {state.logs.slice().reverse().map((log, index) => ( + + {new Date(log.time).toLocaleTimeString('zh-CN', { hour12: false })} · {log.message} + + ))} + +
+
+
+ ); +} diff --git a/web/src/components/accounts/AccountManagementList.tsx b/web/src/components/accounts/AccountManagementList.tsx index fb5e446..41ff3cc 100644 --- a/web/src/components/accounts/AccountManagementList.tsx +++ b/web/src/components/accounts/AccountManagementList.tsx @@ -2,15 +2,16 @@ // 监控账号列表已拆分至 components/accounts/MonitoringAccountList.tsx(card list)。 // 数据走 GET /creator/accounts/monitor-views(画像 + 作品聚合 + 最新画像指标快照)。 // 「启动」= POST /phase-a/accounts/:id/start(恢复账号+启动环境),成功后弹窗展示网关浏览器 -// 登录画面(POST /creator/accounts/:id/login-qr,默认打开 douyin.com/user/self),扫码后核验身份。 +// 登录弹窗(POST /creator/accounts/:id/login-qr),自动核验身份、刷新过期二维码。 import { useCallback, useEffect, useState } from 'react'; import { history } from '@umijs/max'; -import { Alert, App, Button, Card, Flex, Form, Modal, Popconfirm, Select, Space, Table, Tag, Typography } from 'antd'; +import { Alert, App, Avatar, Button, Card, Flex, Form, Modal, Popconfirm, Select, Space, Table, Tag, Typography } from 'antd'; import { PlusOutlined, ReloadOutlined } from '@ant-design/icons'; import type { ColumnsType } from 'antd/es/table'; import { remove, creatorUpdate, creatorGet, creatorAction } from '@/services/api'; import { accountReadiness, conflictMessage, dateTime } from '@/utils/helpers'; import { fixedLeft, fixedRight, tablePagination, tableScroll, useOverflowGrid } from '@/utils/table'; +import AccountLoginModal from '@/components/accounts/AccountLoginModal'; // 仅用 antd 默认组件:Table/Tag/Modal/Popconfirm/Select(tags)。 interface EnvironmentView { @@ -27,6 +28,8 @@ interface Row { id: string; name: string; platform_account_key: string; + avatar_url?: string; + douyin_number?: string; tags?: string[]; runtime_status?: string; environment?: EnvironmentView; @@ -46,15 +49,10 @@ function formatCount(value?: number | null): string { return value >= 10000 ? `${(value / 10000).toFixed(1)}w` : `${value}`; } -interface LoginQR { - content_type: string; - image_base64: string; - qr_detected: boolean; - expires_at: string; -} - export default function AccountManagementList() { const [rows, setRows] = useState([]); + const [environments, setEnvironments] = useState<{ alias: string; runtime_id?: string }[]>([]); + const [environmentTarget, setEnvironmentTarget] = useState(null); const [pending, setPending] = useState(true); const [error, setError] = useState(null); const [tagTarget, setTagTarget] = useState(null); @@ -62,8 +60,6 @@ export default function AccountManagementList() { const [tagBusy, setTagBusy] = useState(false); const [actionBusy, setActionBusy] = useState(''); const [loginTarget, setLoginTarget] = useState(null); - const [loginQR, setLoginQR] = useState(null); - const [loginBusy, setLoginBusy] = useState(false); const { message: messageApi } = App.useApp(); const { vertical, horizontal } = useOverflowGrid(rows.length); @@ -71,7 +67,10 @@ export default function AccountManagementList() { setPending(true); setError(null); try { - const rows = await creatorGet('/creator/accounts/monitor-views'); + const [rows, environments] = await Promise.all([ + creatorGet('/creator/accounts/monitor-views'), creatorGet('/creator/environments'), + ]); + setEnvironments(environments); setRows( (rows ?? []).map((account: any) => ({ ...account, @@ -114,8 +113,6 @@ export default function AccountManagementList() { await creatorAction(`/phase-a/accounts/${encodeURIComponent(row.id)}/start`); messageApi.success('运行环境已启动。'); setLoginTarget(row); - setLoginQR(null); - requestLoginQR(row); } catch (startError) { messageApi.error(conflictMessage(startError, '启动未完成:账号或环境不满足就绪条件,请按不可运行原因处理后重试。')); } finally { @@ -123,37 +120,6 @@ export default function AccountManagementList() { } } - async function requestLoginQR(row?: Row) { - const target = row ?? loginTarget; - if (!target) return; - setLoginBusy(true); - try { - const result = await creatorAction(`/creator/accounts/${encodeURIComponent(target.id)}/login-qr`); - setLoginQR(result); - } catch (qrError) { - messageApi.error(conflictMessage(qrError, '登录画面获取失败;请检查运行环境和网关状态')); - } finally { - setLoginBusy(false); - } - } - - // 扫码后核验:成功视为已登录,刷新列表同步最新账号状态。 - async function verifyLogin() { - if (!loginTarget) return; - setLoginBusy(true); - try { - await creatorAction(`/creator/accounts/${encodeURIComponent(loginTarget.id)}/verify`); - messageApi.success('浏览器身份核验成功,账号已登录。'); - setLoginTarget(null); - setLoginQR(null); - await load(); - } catch (verifyError) { - messageApi.error(conflictMessage(verifyError, '浏览器身份核验失败;请先扫码登录')); - } finally { - setLoginBusy(false); - } - } - async function saveTags() { if (!tagTarget) return; const { tags } = await tagForm.validateFields(); @@ -177,7 +143,7 @@ export default function AccountManagementList() { ...fixedLeft({}), render: (_, account) => (
- history.push(`/accounts/${account.id}`)}>{account.name} + history.push(`/accounts/${account.id}`)}>{account.name} {account.platform_account_key} @@ -290,10 +256,23 @@ export default function AccountManagementList() { + {environmentTarget ? ( + setEnvironmentTarget(null)} onSuccess={() => { setEnvironmentTarget(null); void load(); }} /> + ) : null} + {environments.length > 0 ? ( + + {alias} }, + { title: '状态', render: () => 待登录 }, + { title: '操作', render: (_, env) => }, + ]} /> + + ) : null} {error ? ( 重试} /> ) : null} @@ -323,43 +302,10 @@ export default function AccountManagementList() { - { - if (!loginBusy) { - setLoginTarget(null); - setLoginQR(null); - } - }} - footer={[ - , - , - ]} - > - {loginQR ? ( - - - {loginQR.qr_detected ? '请使用抖音 App 扫描下方二维码。' : '当前画面未检测到二维码,请按画面提示人工完成验证。'} - - 抖音登录二维码或人工验证画面 - - 当前画面有效至 {dateTime(loginQR.expires_at)};完成扫码后请点击“核验浏览器身份”。 - - - ) : ( - 正在获取登录画面… - )} - + {loginTarget ? ( + setLoginTarget(null)} onSuccess={() => { void load(); }} /> + ) : null} ); } diff --git a/web/src/components/accounts/account-login.ts b/web/src/components/accounts/account-login.ts new file mode 100644 index 0000000..c198eb8 --- /dev/null +++ b/web/src/components/accounts/account-login.ts @@ -0,0 +1,101 @@ +import { conflictMessage } from '@/utils/helpers'; + +export interface LoginQR { + content_type: string; + image_base64: string; + qr_detected: boolean; + expires_at: string; +} + +export interface LoginResult { + status: string; + reason?: string; +} + +export interface LoginState { + phase: 'loading' | 'checking' | 'waiting' | 'refreshing' | 'success' | 'error'; + qr?: LoginQR; + error?: string; + logs: { id: number; time: number; message: string }[]; +} + +type LoginRequest = (path: string, options: { method?: string; signal: AbortSignal }) => Promise; + +// A single sequential loop prevents overlapping checks and QR reloads. +// Closing the modal aborts both the current request and the scheduled check. +export function startAccountLogin( + accountID: string, + request: LoginRequest, + onChange: (state: LoginState) => void, + onSuccess: (result: LoginResult) => void, + resource: 'accounts' | 'environments' = 'accounts', +): () => void { + const controller = new AbortController(); + const { signal } = controller; + let timer: ReturnType | undefined; + let state: LoginState = { phase: 'loading', logs: [] }; + let logID = 0; + const path = `/creator/${resource}/${encodeURIComponent(accountID)}`; + const update = (phase: LoginState['phase'], message: string, qr = state.qr, error?: string) => { + if (signal.aborted) return; + state = { phase, qr, error, logs: [...state.logs, { id: ++logID, time: Date.now(), message }].slice(-100) }; + onChange(state); + }; + const pause = (milliseconds: number) => new Promise(resolve => { + const finish = () => { signal.removeEventListener('abort', finish); resolve(); }; + timer = setTimeout(finish, milliseconds); + signal.addEventListener('abort', finish, { once: true }); + }); + const loadQR = async (refresh: boolean) => { + // Remove the expired image immediately, rather than leaving it scannable. + state = { ...state, qr: undefined }; + update(refresh ? 'refreshing' : 'loading', refresh ? '二维码已过期,正在自动刷新。' : '正在获取登录二维码。'); + const qr = await request(`${path}/login-qr`, { + method: 'POST', signal: AbortSignal.any([signal, AbortSignal.timeout(45000)]), + }); + if (signal.aborted) return; + if (!qr?.qr_detected || !qr.image_base64 || !qr.content_type?.startsWith('image/')) { + throw new Error('未获取到可用的登录二维码,请查看浏览器环境。'); + } + const expires = Date.parse(qr.expires_at); + if (!Number.isFinite(expires) || expires <= Date.now()) { + throw new Error('登录二维码有效期无效或已经过期。'); + } + update('waiting', refresh ? '新二维码已加载,等待扫码或手机确认。' : '二维码已加载,等待扫码或手机确认。', qr); + }; + void (async () => { + try { + while (!signal.aborted) { + update(state.qr ? 'checking' : 'loading', state.qr ? '正在核验登录状态。' : '正在检查已有登录会话。'); + const result = await request(`${path}/verify`, { + method: 'POST', signal: AbortSignal.any([signal, AbortSignal.timeout(45000)]), + }); + if (signal.aborted) return; + if (result?.status === 'logged_in') { + update('success', '登录成功,账号身份核验通过。'); + onSuccess(result); + return; + } + if (result?.status !== 'manual_login' || result.reason !== 'awaiting_login') { + throw new Error(`登录核验返回异常状态:${result?.reason || result?.status || '缺少状态'}。`); + } + if (!state.qr) { + await loadQR(false); + continue; + } + if (Date.parse(state.qr.expires_at) <= Date.now()) { + // Verify first: refreshing navigates the browser and could interrupt a completed login. + await loadQR(true); + continue; + } + update('waiting', '尚未完成登录,等待扫码或手机确认。'); + await pause(Math.min(3000, Date.parse(state.qr!.expires_at) - Date.now())); + } + } catch (error) { + if (signal.aborted) return; + const message = conflictMessage(error); + update('error', `登录流程已停止:${message}`, state.qr, message); + } + })(); + return () => { controller.abort(); if (timer !== undefined) clearTimeout(timer); }; +} diff --git a/web/src/pages/accounts/$id/edit.tsx b/web/src/pages/accounts/$id/edit.tsx index 8fb2ac7..c461616 100644 --- a/web/src/pages/accounts/$id/edit.tsx +++ b/web/src/pages/accounts/$id/edit.tsx @@ -4,10 +4,11 @@ // 布局对齐全站模式:实体标识经 usePageInfo 进页头 content 左侧,返回/登录操作按钮经 usePageActions 进页头右侧。 import { useCallback, useEffect, useState } from 'react'; import { history, useParams } from '@umijs/max'; -import { Alert, App, Button, Card, Col, Descriptions, Divider, Flex, Form, Input, Row, Select, Space, Typography } from 'antd'; +import { Alert, App, Button, Card, Col, Descriptions, Flex, Form, Input, Row, Select, Space, Typography } from 'antd'; import { usePageActions, usePageInfo } from '@/components/PageActions'; import FingerprintFields, { fingerprintFormValues, fingerprintPayload } from '@/components/accounts/FingerprintFields'; -import { creatorAction, creatorGet, creatorUpdate, getOne, getList } from '@/services/api'; +import { creatorUpdate, getOne, getList } from '@/services/api'; +import AccountLoginModal from '@/components/accounts/AccountLoginModal'; import { conflictMessage, dateTime } from '@/utils/helpers'; const platformLabelMap: Record = { douyin: '抖音' }; @@ -42,8 +43,8 @@ export default function Page() { const [pending, setPending] = useState(true); const [error, setError] = useState(null); const [busy, setBusy] = useState(false); + const [loginOpen, setLoginOpen] = useState(false); const [fpBusy, setFpBusy] = useState(false); - const [loginQR, setLoginQR] = useState(null); const [form] = Form.useForm(); const [fpForm] = Form.useForm(); const { message: messageApi } = App.useApp(); @@ -54,29 +55,26 @@ export default function Page() { usePageInfo( selected ? ( - + {selected.name || selected.platform_account_key} - {platformLabelMap[selected.platform] || selected.platform} · {selected.platform_account_key} + {platformLabelMap[selected.platform] || selected.platform} · {selected.platform_account_key ? `UID ${selected.platform_account_key}` : 'UID 未绑定'} ) : null, - [selected?.id, selected?.name, selected?.platform], + [selected?.id, selected?.name, selected?.platform, selected?.platform_account_key], ); // 页头 content 行右侧:返回 + 登录操作,不在卡片内重复渲染。 usePageActions( - - , - [selected?.id, selected?.platform], + [selected?.id, selected?.platform, busy, fpBusy, loginOpen], ); const load = useCallback(async () => { @@ -150,45 +148,7 @@ export default function Page() { } } - async function requestLoginQR() { - if (!selected || selected.platform !== 'douyin') return; - setBusy(true); - try { - const result = await creatorAction(`/creator/accounts/${encodeURIComponent(selected.id)}/login-qr`); - setLoginQR(result); - messageApi.success( - result.qr_detected ? '登录二维码已生成,请使用抖音 App 扫码。' : '登录画面已生成;当前未检测到二维码,请按页面提示人工完成验证。', - ); - } catch (qrError) { - messageApi.error(conflictMessage(qrError, '登录二维码获取失败;请检查运行环境和网关状态')); - } finally { - setBusy(false); - } - } - - async function verifyLogin() { - if (!selected) return; - setBusy(true); - try { - const result = await creatorAction(`/creator/accounts/${encodeURIComponent(selected.id)}/verify`); - setProfiles((items) => - items.map((item) => - item.id === selected.id - ? { ...item, login_status: result.status, login_reason: result.reason, login_checked_at: result.checked_at } - : item, - ), - ); - setLoginQR(null); - messageApi.success('浏览器身份核验成功。'); - await load(); - } catch (verifyError) { - messageApi.error(conflictMessage(verifyError, '浏览器身份核验失败;请先在指定环境人工登录')); - } finally { - setBusy(false); - } - } - - if (pending) return ; + if (pending && !selected) return ; if (error && !selected) return (
@@ -202,8 +162,19 @@ export default function Page() { return ( + {loginOpen ? ( + setLoginOpen(false)} onSuccess={() => { void load(); }} /> + ) : null} + + {selected.platform_account_key ? ( + {selected.platform_account_key} + ) : ( + 首次登录核验后自动绑定 + )} + {selected.login_status === 'logged_in' ? '已登录' : '需人工确认'} {selected.login_checked_at ? ` · ${dateTime(selected.login_checked_at)}` : ' · 尚未核验'} @@ -211,22 +182,6 @@ export default function Page() { {selected.password_configured ? '已配置' : '未配置'} - {loginQR ? ( - <> - - - {loginQR.qr_detected ? '请使用抖音 App 扫描下方二维码。' : '当前画面未检测到二维码,请按画面提示人工完成验证。'} - - 抖音登录二维码或人工验证画面 - - 当前画面有效至 {dateTime(loginQR.expires_at)};完成扫码后请点击“核验浏览器身份”。 - - - ) : null}
@@ -284,7 +239,7 @@ export default function Page() { } > - 留空项使用浏览器默认值;seed 由系统按账号自动派生并保证唯一,代理由网络出口统一管理,均不在此配置。运行中的环境保存后会以新指纹自动重启。 + 留空项使用浏览器默认值;seed 在环境创建时分配并保证唯一,账号绑定不会改变它,代理由网络出口统一管理,均不在此配置。运行中的环境保存后会以新指纹自动重启。 diff --git a/web/src/pages/accounts/$id/index.tsx b/web/src/pages/accounts/$id/index.tsx index 4e77f57..4a8ea93 100644 --- a/web/src/pages/accounts/$id/index.tsx +++ b/web/src/pages/accounts/$id/index.tsx @@ -91,16 +91,16 @@ export default function Page() { usePageInfo( account ? ( - + {account.name} {statusTag} - {account.platform_account_key} · {platformLabel(account.platform)} + {account.platform_account_key ? `UID ${account.platform_account_key}` : 'UID 未绑定'} · {platformLabel(account.platform)} ) : null, - [account?.id, account?.name, account?.runtime_status, binding, browsersError], + [account?.id, account?.name, account?.platform_account_key, account?.runtime_status, binding, browsersError], ); const loadAll = useCallback(async () => { @@ -209,6 +209,9 @@ export default function Page() { + + {account.platform_account_key ? {account.platform_account_key} : '首次登录核验后自动绑定'} + {`${account.runtime_status === 'active' ? '启用' : '暂停'} · 版本 ${account.version}`} {account.tags?.join('、') || '无'} diff --git a/web/src/pages/accounts/new.tsx b/web/src/pages/accounts/new.tsx index ffecfe4..2f83066 100644 --- a/web/src/pages/accounts/new.tsx +++ b/web/src/pages/accounts/new.tsx @@ -1,103 +1,44 @@ -// 创建社媒账号:语义对齐 web.archived AccountCreatePage + AccountCreateForm。 -// cookies 非必填:留空代表创建后走扫码登录。创建成功跳编辑页。 -// 折叠区块:指纹浏览器环境(可选,字段复用 FingerprintFields);seed 由后端从账号派生,代理由网络出口管理,均不在表单内。 import { useState } from 'react'; import { history } from '@umijs/max'; -import { Alert, Button, Card, Collapse, Flex, Form, Input, Select, Typography } from 'antd'; -import { create } from '@/services/api'; -import { conflictMessage, platforms } from '@/utils/helpers'; +import { Alert, Button, Card, Collapse, Flex, Form, Typography } from 'antd'; +import { creatorCreate } from '@/services/api'; +import { conflictMessage } from '@/utils/helpers'; import FingerprintFields, { fingerprintPayload, type FingerprintValues } from '@/components/accounts/FingerprintFields'; -interface FormValues { - name: string; - platform: string; - platform_account_key: string; - tags?: string[]; - cookies?: string; - fingerprint?: FingerprintValues; -} - export default function Page() { - const [form] = Form.useForm(); + const [form] = Form.useForm<{ fingerprint?: FingerprintValues }>(); const [busy, setBusy] = useState(false); const [error, setError] = useState(null); - - async function onFinish(values: FormValues) { + async function onFinish(values: { fingerprint?: FingerprintValues }) { setBusy(true); setError(null); try { - const data: Record = { - name: values.name.trim(), - platform: values.platform.trim(), - platform_account_key: values.platform_account_key.trim(), - tags: values.tags ?? [], - }; - // cookies 非必填:留空代表创建后走扫码登录,凭据由后续同步链路补齐 - if (values.cookies?.trim()) data.cookies = values.cookies.trim(); - // 指纹表单:只提交非空项;留空项由后端使用浏览器默认值 - const fingerprint = fingerprintPayload(values.fingerprint); - if (Object.keys(fingerprint).length) data.fingerprint = fingerprint; - const result = await create('accounts', data); - const id = result?.data?.id ?? result?.id; - if (!id) throw new Error('创建账号未返回账号 ID'); - history.replace(`/accounts/${encodeURIComponent(id)}/edit`); - } catch (createError) { - setError(createError); + const environment = await creatorCreate('/creator/environments', { fingerprint: fingerprintPayload(values.fingerprint) }); + if (!environment?.alias) throw new Error('创建环境未返回环境标识'); + history.replace('/accounts'); + } catch (error) { + setError(error); } finally { setBusy(false); } } - return ( - {error ? ( - - ) : null} + {error ? : null} - - - - - - - -