1- from __future__ import annotations
1+ from __future__ import annotations
22
33import locale
44import re
55from contextlib import contextmanager
6+
67from typing import Iterator , Sequence
78
8- from sqlalchemy import URL , Engine , create_engine , select , text
9+ from datetime import datetime , timedelta , timezone
10+ from sqlalchemy import URL , Engine , and_ , or_ , create_engine , select , text
911from sqlalchemy .orm import Session , sessionmaker
1012
1113from config .config import settings
@@ -31,10 +33,10 @@ def _decode_non_utf8_error(exc: UnicodeDecodeError) -> str:
3133 return str (exc )
3234
3335 encodings_to_try = []
34- preferred = locale .getpreferredencoding (False )
36+ preferred = locale .getpreferredencoding (False ) # система
3537 if preferred :
3638 encodings_to_try .append (preferred )
37- encodings_to_try .extend (["cp1251" , "cp866" , "latin1" ])
39+ encodings_to_try .extend (["cp1251" , "cp866" , "latin1" ]) # иные кодировки
3840
3941 seen : set [str ] = set ()
4042 for encoding in encodings_to_try :
@@ -58,7 +60,7 @@ def _build_connection_hint(decoded_message: str, db_name: str) -> str:
5860 ):
5961 return (
6062 f"Database '{ db_name } ' does not exist. "
61- "Run `python main.py --init-only` once or add `--bootstrap` to the run command ."
63+ "Запусти `python main.py --init-only` один раз или используй `--bootstrap` чтобы запустить команду ."
6264 )
6365
6466 if (
@@ -67,18 +69,18 @@ def _build_connection_hint(decoded_message: str, db_name: str) -> str:
6769 ):
6870 return (
6971 "Authentication failed. "
70- "Check DB_HOST/DB_PORT/DB_USER/DB_PASSWORD and PostgreSQL `pg_hba.conf`."
72+ "Проверь DB_HOST/DB_PORT/DB_USER/DB_PASSWORD и PostgreSQL `pg_hba.conf`."
7173 )
7274
7375 if "pg_hba" in lower_message :
7476 return (
75- "Connection rejected by pg_hba.conf. "
76- "Allow the host/user/database combination or use a matching auth method. "
77+ "Подключение отклонено pg_hba.conf. "
78+ "Проверь host/user/database комбинацию или используй соответсвующий метод аутентификации "
7779 )
7880
7981 return (
80- "Connection failed with a non-UTF8 server message. "
81- "Check DB settings and PostgreSQL server logs. "
82+ "Подключение не удалось, сообщение сервера было не в UTF-8. "
83+ "Проверь настройки базы и логи PostgreSQL "
8284 )
8385
8486
@@ -239,8 +241,8 @@ def create_database_if_not_exists(db_name: str) -> None:
239241 # `CREATE DATABASE` cannot be parameterized; reject anything that
240242 # would otherwise require manual quoting/escaping.
241243 raise ValueError (
242- f"Refusing to create database with unsupported name : { db_name !r} . "
243- "Allowed characters: letters, digits, underscore ."
244+ f"Отказано в создании базы данных изза неподобающего имени : { db_name !r} . "
245+ "Разрещенные символы: латиница, цифры, нижнее подчеркивание ."
244246 )
245247
246248 with _connect (engine_for (settings .db_admin_db ), settings .db_admin_db , autocommit = True ) as conn :
@@ -257,6 +259,10 @@ def _create_table(table_attr: str) -> None:
257259
258260def create_search_requests_table () -> None :
259261 _create_table ("search_requests" )
262+ with _connect (_news_engine (), settings .news_db ) as conn :
263+ conn .execute (text (
264+ "ALTER TABLE search_requests ADD COLUMN IF NOT EXISTS is_trial BOOLEAN NOT NULL DEFAULT FALSE"
265+ ))
260266
261267
262268def create_articles_table () -> None :
@@ -265,7 +271,7 @@ def create_articles_table() -> None:
265271
266272def create_user_news_table () -> None :
267273 _create_table ("user_news" )
268-
274+
269275
270276def create_request_stats_table () -> None :
271277 _create_table ("request_stats" )
@@ -336,6 +342,10 @@ def create_request_ai_report_table() -> None:
336342
337343def create_app_users_table () -> None :
338344 _create_table ("app_users" )
345+ with _connect (_news_engine (), settings .news_db ) as conn :
346+ conn .execute (text (
347+ "ALTER TABLE app_users ADD COLUMN IF NOT EXISTS trial_uses Integer not null Default 0"
348+ ))
339349
340350
341351def create_news_tables () -> None :
@@ -344,7 +354,16 @@ def create_news_tables() -> None:
344354
345355def create_users_keys_table () -> None :
346356 _create_table ("users_keys" )
347-
357+ migrations = [
358+ "ALTER TABLE users_keys ADD COLUMN IF NOT EXISTS status TEXT NOT NULL DEFAULT 'pending_validation'" ,
359+ "ALTER TABLE users_keys ADD COLUMN IF NOT EXISTS validation_error TEXT" ,
360+ "ALTER TABLE users_keys ADD COLUMN IF NOT EXISTS validated_at TIMESTAMPTZ" ,
361+ ]
362+ add_status_check = """
363+ ALTER TABLE users_keys DROP CONSTRAINT IF EXISTS users_keys_status_check;
364+ ALTER TABLE users_keys ADD CONSTRAINT users_keys_status_check
365+ CHECK (status IN ('pending_validation', 'validating', 'valid', 'invalid', 'exhausted'));
366+ """
348367 trigger_function = """
349368 CREATE OR REPLACE FUNCTION set_users_keys_updated_at()
350369 RETURNS TRIGGER AS $$
@@ -353,11 +372,10 @@ def create_users_keys_table() -> None:
353372 RETURN NEW;
354373 END;
355374 $$ LANGUAGE plpgsql;
356- """
357-
375+ """
358376 trigger = """
359377 DO $$
360- BEGIN
378+ BEGIN
361379 IF NOT EXISTS (
362380 SELECT 1 FROM pg_trigger WHERE tgname = 'trg_users_keys_updated_at'
363381 ) THEN
@@ -367,41 +385,71 @@ def create_users_keys_table() -> None:
367385 EXECUTE FUNCTION set_users_keys_updated_at();
368386 END IF;
369387 END $$;
370- """
371-
388+ """
372389 with _connect (_news_engine (), settings .news_db ) as conn :
390+ for m in migrations :
391+ conn .execute (text (m ))
392+ conn .execute (text (add_status_check ))
393+ conn .execute (text (
394+ "CREATE INDEX IF NOT EXISTS users_keys_status_idx ON users_keys (status)"
395+ ))
373396 conn .execute (text (trigger_function ))
374397 conn .execute (text (trigger ))
375398
376399
377400def claim_next_search_request () -> dict | None :
378401 # CTE-based atomic claim: SELECT ... FOR UPDATE SKIP LOCKED + UPDATE in
379402 # one statement so multiple workers can run safely against the queue.
380- query = text (
381- """
382- WITH next_request AS (
383- SELECT id
384- FROM search_requests
385- WHERE status = 'queued'
386- ORDER BY created_at
387- FOR UPDATE SKIP LOCKED
388- LIMIT 1
389- )
390- UPDATE search_requests AS sr
391- SET
392- status = 'running',
393- started_at = NOW(),
394- error_text = NULL
395- FROM next_request
396- WHERE sr.id = next_request.id
397- RETURNING sr.id, sr.user_id, sr.keyword, sr.language, sr.limit_count, sr.page_size
398- """
399- )
400-
401- with _connect (_news_engine (), settings .news_db ) as conn :
402- row = conn .execute (query ).mappings ().first ()
403- return dict (row ) if row else None
404-
403+ with get_session () as session :
404+ request = session .execute (
405+ select (SearchRequest ).where (SearchRequest .status == "queued" ).order_by (SearchRequest .created_at ).limit (1 ).with_for_update (skip_locked = True )
406+ ).scalar_one_or_none ()
407+ if request is None :
408+ return None
409+ request .status = "running"
410+ request .started_at = datetime .now (timezone .utc )
411+ request .error_text = None
412+ session .flush ()
413+
414+ return {
415+ "id" : request .id ,
416+ "user_id" : request .user_id ,
417+ "keyword" : request .keyword ,
418+ "language" : request .language ,
419+ "limit_count" : request .limit_count ,
420+ "page_size" : request .page_size ,
421+ "is_trial" : request .is_trial
422+ }
423+
424+
425+ def claim_pending_key_validation () -> dict | None :
426+ # подхватывает и pending и зависшие
427+ # в случае падения воркера триггер updated_at покажет что просшло > 10 мин
428+ stale_threshold = datetime .now (timezone .utc ) - timedelta (minutes = 10 )
429+ with get_session () as session :
430+ key = session .execute (select (UsersKeys ).where (or_ (UsersKeys .status == "pending_validation" , and_ (UsersKeys .status == "validating" , UsersKeys .updated_at < stale_threshold ,)))
431+ .order_by (UsersKeys .id ).limit (1 ).with_for_update (skip_locked = True )).scalar_one_or_none ()
432+ if key is None :
433+ return None
434+ key .status = "validating"
435+ session .flush ()
436+
437+ return {
438+ "id" : key .id ,
439+ "user_id" : key .user_id ,
440+ "encrypted_key" : key .encrypted_key ,
441+ "iv" : key .iv ,
442+ "auth_tag" : key .auth_tag ,
443+ }
444+ def mark_key_validation_result (key_id : int , status : str , error : str | None ) -> None :
445+ with get_session () as session :
446+ key = session .get (UsersKeys , key_id )
447+ if key is None :
448+ return
449+ key .status = status
450+ key .validation_error = error
451+ if status != "pending_validation" :
452+ key .validated_at = datetime .now (timezone .utc )
405453
406454def search_request_exists (search_request_id : int ) -> bool :
407455 with get_session () as session :
@@ -477,4 +525,7 @@ def fetch_articles_for_search_request(search_request_id: int) -> list[dict]:
477525 "search_request_belongs_to_user" ,
478526 "search_request_exists" ,
479527 "table_exists" ,
528+ "claim_panding_key_validation" ,
529+ "mark_key_validation_result"
480530]
531+
0 commit comments