Skip to content

Commit 3649a9d

Browse files
committed
Update: Генерация отклика через threading
1 parent 2348488 commit 3649a9d

3 files changed

Lines changed: 150 additions & 54 deletions

File tree

.gitignore

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,35 @@
1+
# Python
12
__pycache__/
23
*.pyc
4+
*.pyo
5+
*.pyd
6+
*.egg-info/
7+
dist/
8+
build/
9+
.venv/
10+
venv/
11+
12+
# Env & secrets
313
.env
14+
.env.*
15+
!.env.example
16+
17+
# Project data
418
data/
19+
*.db
20+
*.sqlite
21+
*.sqlite3
22+
23+
# RuFlo / claude-flow runtime
24+
.claude-flow/
25+
.swarm/
26+
.claude/memory.db
27+
28+
# Logs & temp
29+
*.log
30+
tmp/
31+
temp/
32+
33+
# OS
34+
.DS_Store
35+
Thumbs.db

kwork_parser/app.py

Lines changed: 78 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
import html
44
import logging
5+
import threading
56
import time
67
from datetime import datetime
78

@@ -22,6 +23,9 @@
2223
logger = logging.getLogger(__name__)
2324

2425

26+
_MAX_CONCURRENT_DRAFTS = 3
27+
28+
2529
class Application:
2630
def __init__(self, settings: Settings) -> None:
2731
self.settings = settings
@@ -31,6 +35,7 @@ def __init__(self, settings: Settings) -> None:
3135
self.ai_scorer = OpenRouterScorer(settings) if settings.ai_enabled else None
3236
self.response_draft_generator = ResponseDraftService(settings) if settings.response_draft_enabled else None
3337
self.notifier = TelegramNotifier(settings)
38+
self._draft_semaphore = threading.Semaphore(_MAX_CONCURRENT_DRAFTS)
3439

3540
def run_once(self) -> int:
3641
self._sync_telegram_feedback()
@@ -295,26 +300,41 @@ def _handle_draft_action(self, action: TelegramDraftAction) -> None:
295300
self.notifier.answer_feedback(action.callback_query_id, "Для этого заказа пока не хватает данных на демо")
296301
return
297302
self.notifier.answer_feedback(action.callback_query_id, "Готовлю демо...")
298-
generated, error_text = self._send_demo_project(
299-
candidate.project,
300-
candidate.rule_result,
301-
candidate.ai_result,
302-
demo_summary=draft.demo_summary,
303-
chat_id=action.chat_id,
304-
)
305-
if not generated:
306-
try:
307-
self.notifier.send_demo_status(candidate.project, error_text, chat_id=action.chat_id)
308-
except Exception:
309-
logger.warning(
310-
"Failed to send demo error status for project %s",
311-
action.project_id,
312-
exc_info=True,
313-
)
303+
threading.Thread(
304+
target=self._run_demo_in_background,
305+
args=(action, candidate, draft),
306+
daemon=True,
307+
).start()
314308
return
315309

316310
self.notifier.answer_feedback(action.callback_query_id, "Генерирую отклик...")
317311
variant = "default" if action.action in {"generate", "regenerate"} else action.action
312+
threading.Thread(
313+
target=self._run_draft_in_background,
314+
args=(action, candidate, variant),
315+
daemon=True,
316+
).start()
317+
except Exception:
318+
logger.warning("Draft action handling failed for project %s", action.project_id, exc_info=True)
319+
320+
def _run_draft_in_background(
321+
self,
322+
action: TelegramDraftAction,
323+
candidate: object,
324+
variant: str,
325+
) -> None:
326+
if not self._draft_semaphore.acquire(blocking=False):
327+
logger.warning("Draft semaphore full, skipping generation for project %s", action.project_id)
328+
try:
329+
self.notifier.send_demo_status(
330+
candidate.project,
331+
"Слишком много одновременных запросов, попробуй чуть позже",
332+
chat_id=action.chat_id,
333+
)
334+
except Exception:
335+
pass
336+
return
337+
try:
318338
self._send_response_draft(
319339
candidate.project,
320340
candidate.rule_result,
@@ -323,7 +343,48 @@ def _handle_draft_action(self, action: TelegramDraftAction) -> None:
323343
chat_id=action.chat_id,
324344
)
325345
except Exception:
326-
logger.warning("Draft action handling failed for project %s", action.project_id, exc_info=True)
346+
logger.warning("Background draft generation failed for project %s", action.project_id, exc_info=True)
347+
finally:
348+
self._draft_semaphore.release()
349+
350+
def _run_demo_in_background(
351+
self,
352+
action: TelegramDraftAction,
353+
candidate: object,
354+
draft: object,
355+
) -> None:
356+
if not self._draft_semaphore.acquire(blocking=False):
357+
logger.warning("Draft semaphore full, skipping demo for project %s", action.project_id)
358+
try:
359+
self.notifier.send_demo_status(
360+
candidate.project,
361+
"Слишком много одновременных запросов, попробуй чуть позже",
362+
chat_id=action.chat_id,
363+
)
364+
except Exception:
365+
pass
366+
return
367+
try:
368+
generated, error_text = self._send_demo_project(
369+
candidate.project,
370+
candidate.rule_result,
371+
candidate.ai_result,
372+
demo_summary=draft.demo_summary,
373+
chat_id=action.chat_id,
374+
)
375+
if not generated:
376+
try:
377+
self.notifier.send_demo_status(candidate.project, error_text, chat_id=action.chat_id)
378+
except Exception:
379+
logger.warning(
380+
"Failed to send demo error status for project %s",
381+
action.project_id,
382+
exc_info=True,
383+
)
384+
except Exception:
385+
logger.warning("Background demo generation failed for project %s", action.project_id, exc_info=True)
386+
finally:
387+
self._draft_semaphore.release()
327388

328389
def _format_health_message(self) -> str:
329390
snapshot = self.storage.get_health_snapshot()

kwork_parser/storage.py

Lines changed: 41 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22

33
import json
44
import sqlite3
5+
import threading
56
from dataclasses import dataclass
67
from pathlib import Path
78

@@ -51,8 +52,9 @@ class Storage:
5152
def __init__(self, path: Path) -> None:
5253
self.path = path
5354
self.path.parent.mkdir(parents=True, exist_ok=True)
54-
self.connection = sqlite3.connect(self.path)
55+
self.connection = sqlite3.connect(self.path, check_same_thread=False)
5556
self.connection.row_factory = sqlite3.Row
57+
self._lock = threading.Lock()
5658
self._init_schema()
5759

5860
def _init_schema(self) -> None:
@@ -372,32 +374,33 @@ def get_health_snapshot(self) -> HealthSnapshot:
372374
)
373375

374376
def save_response_draft(self, draft: ResponseDraft) -> None:
375-
self.connection.execute(
376-
"""
377-
INSERT INTO response_drafts (
378-
project_id, text, variant, demo_available, demo_summary,
379-
demo_path, demo_archive_path, status, updated_at
377+
with self._lock:
378+
self.connection.execute(
379+
"""
380+
INSERT INTO response_drafts (
381+
project_id, text, variant, demo_available, demo_summary,
382+
demo_path, demo_archive_path, status, updated_at
383+
)
384+
VALUES (?, ?, ?, ?, ?, NULL, NULL, 'draft', CURRENT_TIMESTAMP)
385+
ON CONFLICT(project_id) DO UPDATE SET
386+
text = excluded.text,
387+
variant = excluded.variant,
388+
demo_available = excluded.demo_available,
389+
demo_summary = excluded.demo_summary,
390+
demo_path = NULL,
391+
demo_archive_path = NULL,
392+
status = 'draft',
393+
updated_at = CURRENT_TIMESTAMP
394+
""",
395+
(
396+
draft.project_id,
397+
draft.text,
398+
draft.variant,
399+
1 if draft.demo_available else 0,
400+
draft.demo_summary or None,
401+
),
380402
)
381-
VALUES (?, ?, ?, ?, ?, NULL, NULL, 'draft', CURRENT_TIMESTAMP)
382-
ON CONFLICT(project_id) DO UPDATE SET
383-
text = excluded.text,
384-
variant = excluded.variant,
385-
demo_available = excluded.demo_available,
386-
demo_summary = excluded.demo_summary,
387-
demo_path = NULL,
388-
demo_archive_path = NULL,
389-
status = 'draft',
390-
updated_at = CURRENT_TIMESTAMP
391-
""",
392-
(
393-
draft.project_id,
394-
draft.text,
395-
draft.variant,
396-
1 if draft.demo_available else 0,
397-
draft.demo_summary or None,
398-
),
399-
)
400-
self.connection.commit()
403+
self.connection.commit()
401404

402405
def get_response_draft(self, project_id: int) -> ResponseDraft | None:
403406
row = self.connection.execute(
@@ -421,17 +424,18 @@ def get_response_draft(self, project_id: int) -> ResponseDraft | None:
421424
)
422425

423426
def save_demo_project_artifacts(self, project_id: int, demo_path: str, demo_archive_path: str) -> None:
424-
self.connection.execute(
425-
"""
426-
UPDATE response_drafts
427-
SET demo_path = ?,
428-
demo_archive_path = ?,
429-
updated_at = CURRENT_TIMESTAMP
430-
WHERE project_id = ?
431-
""",
432-
(demo_path, demo_archive_path, project_id),
433-
)
434-
self.connection.commit()
427+
with self._lock:
428+
self.connection.execute(
429+
"""
430+
UPDATE response_drafts
431+
SET demo_path = ?,
432+
demo_archive_path = ?,
433+
updated_at = CURRENT_TIMESTAMP
434+
WHERE project_id = ?
435+
""",
436+
(demo_path, demo_archive_path, project_id),
437+
)
438+
self.connection.commit()
435439

436440
def mark_response_draft_sent_manually(self, project_id: int) -> None:
437441
self.connection.execute(

0 commit comments

Comments
 (0)