Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
ddc9cba
docs: design public account id extractor
yukikojo Jul 24, 2026
5baeec6
docs: design keyword crawl runner
yukikojo Jul 27, 2026
1e20b58
docs: plan keyword crawl runner implementation
yukikojo Jul 27, 2026
2d64b9e
chore: ignore local worktrees
yukikojo Jul 27, 2026
c69cc70
feat: expose crawler task completion status
yukikojo Jul 27, 2026
74d7f4a
test: close mocked crawler output coroutine
yukikojo Jul 27, 2026
7de6bcb
fix: isolate crawler reader lifecycle
yukikojo Jul 27, 2026
6a4ba2e
fix: confirm crawler termination before cleanup
yukikojo Jul 27, 2026
f0470b0
docs: design JSON keyword queue upgrade
yukikojo Jul 27, 2026
30ac532
docs: plan JSON keyword queue upgrade
yukikojo Jul 27, 2026
a51febf
fix: prefer discovered CDP websocket
yukikojo Jul 28, 2026
429f5f7
fix: harden CDP endpoint discovery
yukikojo Jul 28, 2026
582b855
fix: validate CDP websocket endpoints
yukikojo Jul 28, 2026
416825e
fix: reject whitespace in CDP websocket URL
yukikojo Jul 28, 2026
f553071
fix: wait for manual XHS CAPTCHA
yukikojo Jul 28, 2026
3cc5590
fix: refresh XHS state after CAPTCHA
yukikojo Jul 28, 2026
34712fa
fix: present XHS verification page
yukikojo Jul 28, 2026
5f85ad4
fix: reuse XHS search context for CAPTCHA
yukikojo Jul 28, 2026
ddcd292
docs: design XHS needs-review partial saves
yukikojo Jul 29, 2026
1ee8995
docs: design public account creator pipeline
yukikojo Jul 31, 2026
a8ce533
docs: capture creator ids during initial crawl
yukikojo Jul 31, 2026
d64039f
docs: plan XHS creator id pipeline
yukikojo Jul 31, 2026
0f07d40
feat: capture XHS creator ids during initial search
yukikojo Jul 31, 2026
a748630
feat: expose XHS creator id capture switch
yukikojo Jul 31, 2026
3c3e4fb
fix: record invalid creator id captures safely
yukikojo Jul 31, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,7 @@ cython_debug/
/temp_image/
/browser_data/
/data/
/data/xhs/private/

*/.DS_Store
.vscode
Expand All @@ -183,4 +184,5 @@ agent_zone
debug_tools

database/*.db
.omx/
.omx/
.worktrees/
10 changes: 7 additions & 3 deletions api/routers/crawler.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,14 +27,18 @@
@router.post("/start")
async def start_crawler(request: CrawlerStartRequest):
"""Start crawler task"""
success = await crawler_manager.start(request)
if not success:
task_id = await crawler_manager.start(request)
if not task_id:
# Handle concurrent/duplicate requests: if process is already running, return 400 instead of 500
if crawler_manager.process and crawler_manager.process.poll() is None:
raise HTTPException(status_code=400, detail="Crawler is already running")
raise HTTPException(status_code=500, detail="Failed to start crawler")

return {"status": "ok", "message": "Crawler started successfully"}
return {
"status": "ok",
"message": "Crawler started successfully",
"task_id": task_id,
}


@router.post("/stop")
Expand Down
4 changes: 4 additions & 0 deletions api/schemas/crawler.py
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ class CrawlerStartRequest(BaseModel):
start_page: int = 1
enable_comments: bool = True
enable_sub_comments: bool = False
capture_creator_ids: bool = False
save_option: SaveDataOptionEnum = SaveDataOptionEnum.JSONL
cookies: str = ""
headless: bool = False
Expand All @@ -85,6 +86,9 @@ class CrawlerStatusResponse(BaseModel):
crawler_type: Optional[str] = None
started_at: Optional[str] = None
error_message: Optional[str] = None
task_id: Optional[str] = None
last_exit_code: Optional[int] = None
finished_at: Optional[str] = None


class LogEntry(BaseModel):
Expand Down
128 changes: 94 additions & 34 deletions api/services/crawler_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
from typing import Optional, List
from datetime import datetime
from pathlib import Path
from uuid import uuid4

from ..schemas import CrawlerStartRequest, LogEntry

Expand All @@ -36,6 +37,10 @@ def __init__(self):
self.status = "idle"
self.started_at: Optional[datetime] = None
self.current_config: Optional[CrawlerStartRequest] = None
self.task_id: Optional[str] = None
self.last_exit_code: Optional[int] = None
self.finished_at: Optional[datetime] = None
self.error_message: Optional[str] = None
self._log_id = 0
self._logs: List[LogEntry] = []
self._read_task: Optional[asyncio.Task] = None
Expand Down Expand Up @@ -90,11 +95,15 @@ def _parse_log_level(self, line: str) -> str:
return "debug"
return "info"

async def start(self, config: CrawlerStartRequest) -> bool:
"""Start crawler process"""
async def start(self, config: CrawlerStartRequest) -> Optional[str]:
"""Start crawler process and return its task id."""
async with self._lock:
if self.process and self.process.poll() is None:
return False
if (
self._read_task and not self._read_task.done()
) or (
self.process and self.process.poll() is None
):
return None

# Clear old logs
self._logs = []
Expand Down Expand Up @@ -130,6 +139,10 @@ async def start(self, config: CrawlerStartRequest) -> bool:
env={**os.environ, "PYTHONUNBUFFERED": "1"}
)

self.task_id = str(uuid4())
self.last_exit_code = None
self.finished_at = None
self.error_message = None
self.status = "running"
self.started_at = datetime.now()
self.current_config = config
Expand All @@ -143,52 +156,76 @@ async def start(self, config: CrawlerStartRequest) -> bool:
# Start log reading task
self._read_task = asyncio.create_task(self._read_output())

return True
return self.task_id
except Exception as e:
self.status = "error"
entry = self._create_log_entry(f"Failed to start crawler: {str(e)}", "error")
await self._push_log(entry)
return False
return None

async def stop(self) -> bool:
"""Stop crawler process"""
async with self._lock:
if not self.process or self.process.poll() is not None:
process = self.process
if not process or process.poll() is not None:
return False

self.status = "stopping"
entry = self._create_log_entry("Sending SIGTERM to crawler process...", "warning")
await self._push_log(entry)
errors = []

try:
self.process.send_signal(signal.SIGTERM)

process.send_signal(signal.SIGTERM)
except Exception as e:
errors.append(str(e))
else:
# Wait for graceful exit (up to 15 seconds)
for _ in range(30):
if self.process.poll() is not None:
if process.poll() is not None:
break
await asyncio.sleep(0.5)

# If still not exited, force kill
if self.process.poll() is None:
entry = self._create_log_entry("Process not responding, sending SIGKILL...", "warning")
await self._push_log(entry)
self.process.kill()

entry = self._create_log_entry("Crawler process terminated", "info")
# If still not exited, force kill and wait for confirmed exit.
if process.poll() is None:
entry = self._create_log_entry("Process not responding, sending SIGKILL...", "warning")
await self._push_log(entry)

except Exception as e:
entry = self._create_log_entry(f"Error stopping crawler: {str(e)}", "error")
try:
process.kill()
await asyncio.to_thread(process.wait, timeout=5)
except Exception as e:
errors.append(str(e))

exit_code = process.poll()
if exit_code is None:
message = "Error stopping crawler: " + "; ".join(
errors or ["process did not exit"]
)
self.status = "error"
self.error_message = message
entry = self._create_log_entry(message, "error")
await self._push_log(entry)
return False

entry = self._create_log_entry("Crawler process terminated", "info")
await self._push_log(entry)
self.status = "idle"
self.current_config = None

# Cancel log reading task
if self._read_task:
self._read_task.cancel()
self._read_task = None
self.last_exit_code = exit_code
self.finished_at = datetime.now()
self.error_message = None

# Cancel and settle log reader before allowing a new start.
read_task = self._read_task
if read_task:
read_task.cancel()
try:
await read_task
except asyncio.CancelledError:
pass
finally:
if self._read_task is read_task:
self._read_task = None

return True

Expand All @@ -199,7 +236,10 @@ def get_status(self) -> dict:
"platform": self.current_config.platform.value if self.current_config else None,
"crawler_type": self.current_config.crawler_type.value if self.current_config else None,
"started_at": self.started_at.isoformat() if self.started_at else None,
"error_message": None
"error_message": self.error_message,
"task_id": self.task_id,
"last_exit_code": self.last_exit_code,
"finished_at": self.finished_at.isoformat() if self.finished_at else None,
}

def _build_command(self, config: CrawlerStartRequest) -> list:
Expand All @@ -225,6 +265,9 @@ def _build_command(self, config: CrawlerStartRequest) -> list:
cmd.extend(["--get_comment", "true" if config.enable_comments else "false"])
cmd.extend(["--get_sub_comment", "true" if config.enable_sub_comments else "false"])

if config.capture_creator_ids:
cmd.extend(["--capture_creator_ids", "true"])

if config.max_notes_count is not None:
cmd.extend(["--crawler_max_notes_count", str(config.max_notes_count)])

Expand All @@ -241,12 +284,14 @@ def _build_command(self, config: CrawlerStartRequest) -> list:
async def _read_output(self):
"""Asynchronously read process output"""
loop = asyncio.get_event_loop()
process = self.process
task_id = self.task_id

try:
while self.process and self.process.poll() is None:
while process and process.poll() is None:
# Read a line in thread pool
line = await loop.run_in_executor(
None, self.process.stdout.readline
None, process.stdout.readline
)
if line:
line = line.strip()
Expand All @@ -256,9 +301,9 @@ async def _read_output(self):
await self._push_log(entry)

# Read remaining output
if self.process and self.process.stdout:
if process and process.stdout:
remaining = await loop.run_in_executor(
None, self.process.stdout.read
None, process.stdout.read
)
if remaining:
for line in remaining.strip().split('\n'):
Expand All @@ -268,20 +313,35 @@ async def _read_output(self):
await self._push_log(entry)

# Process ended
if self.status == "running":
exit_code = self.process.returncode if self.process else -1
if (
self.process is process
and self.task_id == task_id
and self.status == "running"
):
exit_code = process.returncode if process else -1
self.last_exit_code = exit_code
self.finished_at = datetime.now()
self.error_message = None
if exit_code == 0:
entry = self._create_log_entry("Crawler completed successfully", "success")
else:
entry = self._create_log_entry(f"Crawler exited with code: {exit_code}", "warning")
await self._push_log(entry)
self.status = "idle"
self.current_config = None

except asyncio.CancelledError:
pass
except Exception as e:
entry = self._create_log_entry(f"Error reading output: {str(e)}", "error")
await self._push_log(entry)
if self.process is process and self.task_id == task_id:
message = f"Error reading output: {str(e)}"
returncode = process.returncode if process else None
self.last_exit_code = returncode if isinstance(returncode, int) else -1
self.finished_at = datetime.now()
self.error_message = message
self.status = "error"
entry = self._create_log_entry(message, "error")
await self._push_log(entry)


# Global singleton
Expand Down
11 changes: 11 additions & 0 deletions cmd_arg/arg.py
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,15 @@ def main(
show_default=True,
),
] = str(config.ENABLE_GET_SUB_COMMENTS),
capture_creator_ids: Annotated[
str,
typer.Option(
"--capture_creator_ids",
help="Whether to capture creator IDs from XHS search results, supports yes/true/t/y/1 or no/false/f/n/0",
rich_help_panel="Basic Configuration",
show_default=True,
),
] = str(config.ENABLE_XHS_CREATOR_ID_CAPTURE),
headless: Annotated[
str,
typer.Option(
Expand Down Expand Up @@ -337,6 +346,7 @@ def main(

enable_comment = _to_bool(get_comment)
enable_sub_comment = _to_bool(get_sub_comment)
enable_creator_id_capture = _to_bool(capture_creator_ids)
enable_headless = _to_bool(headless)
enable_ip_proxy_value = _to_bool(enable_ip_proxy)
init_db_value = init_db.value if init_db else None
Expand All @@ -353,6 +363,7 @@ def main(
config.KEYWORDS = keywords
config.ENABLE_GET_COMMENTS = enable_comment
config.ENABLE_GET_SUB_COMMENTS = enable_sub_comment
config.ENABLE_XHS_CREATOR_ID_CAPTURE = enable_creator_id_capture
config.HEADLESS = enable_headless
config.CDP_HEADLESS = enable_headless
config.SAVE_DATA_OPTION = save_data_option.value
Expand Down
4 changes: 4 additions & 0 deletions config/base_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,10 @@
# Data saving path, if not specified by default, it will be saved to the data folder.
SAVE_DATA_PATH = ""

# Opt-in local sidecar for XHS creator IDs captured during keyword search.
ENABLE_XHS_CREATOR_ID_CAPTURE = False
XHS_CREATOR_ID_CAPTURE_DIR = ""

# Browser file configuration cached by the user's browser
USER_DATA_DIR = "%s_user_data_dir" # %s will be replaced by platform name

Expand Down
Loading