-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathmain.py
More file actions
2814 lines (2486 loc) · 163 KB
/
Copy pathmain.py
File metadata and controls
2814 lines (2486 loc) · 163 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
import sys
# Кодировку stdio выставляем ДО первого импорта: config.py печатает статус с
# эмодзи прямо при импорте, и когда stdout не UTF-8 — запуск мимо start.bat,
# перенаправление в лог-файл, супервизор, планировщик задач — этот print
# падает с UnicodeEncodeError ещё до старта бота. В логе не остаётся ничего:
# логирование к тому моменту не настроено. Тот же класс отказа, что и с
# потерянным DENTAL_KEYWORDS — бот просто не поднимается.
for _stream in (sys.stdout, sys.stderr):
try:
_stream.reconfigure(encoding="utf-8", errors="replace")
except Exception:
pass
import logging
from telethon import TelegramClient, events
from telethon import utils as telethon_utils
import config
import runtime_guard
runtime_guard.configure_logging()
logger = logging.getLogger(__name__)
import vision
import os
import re
import time
import asyncio
import json
from collections import deque
import database
import assistant
from datetime import datetime
from datetime import timedelta
from datetime import timezone
import summarizer
from media_tools import (MEDIA_TEMP_DIR as media_temp_dir, clinical_media_kind,
extract_first_frame_async, image_document)
try:
import psutil
except Exception:
psutil = None
PROCESSED_MSG_IDS = []
def _env_int(name, default):
try:
return int(os.getenv(name, str(default)))
except ValueError:
return default
MY_ID = 7716348189
# Числовой id бота как запасной вариант. Основной источник —
# assistant.BOT_ID, но он появляется только после init_assistant: до этого
# момента (и если get_me не прошёл) свои сообщения всё равно надо уметь
# опознавать, иначе бот отвечает на собственный дайджест.
FALLBACK_BOT_ID = _env_int("STOMCHAT_BOT_ID", 7971556097)
# Имя, по которому бота зовут в группе, пока assistant.BOT_USERNAME не
# отрезолвился. До этого литерал был единственным работающим способом:
# соседняя проверка f"@{assistant.BOT_ID}" сравнивала текст с числовым id и
# не срабатывала никогда.
FALLBACK_BOT_USERNAME = os.getenv("STOMCHAT_BOT_USERNAME", "stomchat_bot").lstrip("@").lower()
HEALTH_CHECK_INTERVAL_SECONDS = 300
HEALTH_FAILURE_LIMIT = 3
SCHEDULER_STATE_PATH = "bot_state.json"
SUMMARY_STATUS_CHECK_SECONDS = 60
# Терпение сторожа сводки обязано БЫТЬ БОЛЬШЕ самого долгого законного шага
# конвейера, иначе сторож стреляет в живую работу.
#
# Здесь стояло 1800, а генерации сводки разрешено GEMINI_GENERATION_TIMEOUT_SECONDS
# = 2100. Разбор путей: на одну попытку к провайдеру уходит timeout/3 = 700 с, и
# статус пишется ПЕРЕД каждой попыткой, поэтому при живом дочернем процессе
# разрыв между записями не превышает 700 с и сторож молчит. Но если ребёнок жив
# и молчит до первой попытки — завис на импорте или на DNS при создании
# клиента, — записей нет вовсе. Тогда на 1800 с сторож убивал процесс, хотя
# родитель отпустил бы вызов по своему таймауту на 2100 с и обработал отказ
# штатно: записал бы ошибку, снял флаг и вернул None. Вместо этого терялся
# дайджест И происходил перезапуск.
#
# Значение выведено из бюджета генерации, а не выбрано отдельным числом:
# разъехавшись, они вернут то же противоречие. Запас в 600 с покрывает
# публикацию (60 с), отправку по целям (90 с на цель) и закрепление (30 с).
# Сторож остаётся страховкой от настоящего зависания — просто перестаёт быть
# быстрее собственных таймаутов конвейера.
SUMMARY_STALE_SECONDS = summarizer.GEMINI_GENERATION_TIMEOUT_SECONDS + 600
START_TIMEOUT_SECONDS = 120
# Догоняющая синхронизация упирается в размер долга, а не в скорость сети.
# Порядок величины: по локальному снимку базы в чате около 228 сообщений в
# сутки, то есть неделя простоя — это порядка 1600 реплик, месяц — порядка 8000.
# Снимок локальный и может отставать от боевого, поэтому это оценка масштаба, а
# не факт о бое; сам вывод от неё не зависит.
#
# В прежние 300 с восемь тысяч укладывались только при 28 сообщениях в секунду.
# Не уложившись, синхронизация валила ВЕСЬ подъём: start.bat поднимал процесс
# заново, и так по кругу. Сообщения сохраняются по ходу, поэтому каждый заход
# продвигался, но снаружи это выглядит как «бот не запускается».
#
# Держать долгую синхронизацию безопасно: внутри цикла каждые 25 сообщений
# пишется heartbeat, и сторож видит процесс живым.
SYNC_HISTORY_TIMEOUT_SECONDS = 900
TELEGRAM_REQUEST_TIMEOUT_SECONDS = 60
# Автор сообщения — необязательное обогащение: при отказе в базу пишется
# "Unknown", и таких строк уже 917. Держать здесь общий таймаут в 60 с нельзя:
# пяти зависших запросов хватало, чтобы съесть весь бюджет подъёма.
SYNC_SENDER_TIMEOUT_SECONDS = 10
MEDIA_DOWNLOAD_TIMEOUT_SECONDS = 120
# Внешний потолок разбора снимка ОБЯЗАН вмещать внутренний бюджет каскада зрения,
# иначе резервные модели не пробуются никогда.
#
# Замер по константам vision.py и живому набору ключей: пул из трёх моделей
# перебирал ВСЕ ключи своего провайдера — 2 x 10 + 1 x 7 = 27 попыток по 33 с
# (запрос 30 плюс пауза троттлинга 3), до 891 с на один снимок. Потолок стоял
# 180 с, то есть в пять раз меньше: резервные модели каскада не пробовались
# никогда, а снимок после срабатывания внешнего таймаута получал в базе
# описание "-" — отметку «разобрано», — и второго шанса у него не было:
# get_pending_media_message_ids его больше не возвращает. Рентген коллеги
# терялся молча и навсегда.
#
# Поднимать потолок до 951 с было бы ХУЖЕ: воркер один, очередь на 128 снимков,
# и одно зависшее фото остановило бы разбор на четверть часа. Поэтому ограничен
# сам каскад — не более VISION_KEYS_PER_MODEL ключей на модель, — а потолок
# ВЫВЕДЕН из его фактического бюджета, а не задан отдельным числом: два
# независимых числа уже трижды за сессию разъезжались (test_budget_nesting.py).
# Запас 60 с — на подготовку снимка (до 45 с) и накладные.
MEDIA_ANALYSIS_TIMEOUT_SECONDS = max(
180,
_env_int("STOMCHAT_MEDIA_ANALYSIS_TIMEOUT_SECONDS", 0)
or (vision.vision_cascade_budget_seconds() + 60),
)
MEDIA_FRAME_TIMEOUT_SECONDS = 60
MEDIA_WORKER_COUNT = max(1, _env_int("STOMCHAT_MEDIA_WORKERS", 1))
MEDIA_QUEUE_MAX_SIZE = max(MEDIA_WORKER_COUNT, _env_int("STOMCHAT_MEDIA_QUEUE_MAX", 128))
# Каталог общий на проект, объявлен в media_tools: его же используют
# уборщик, голосовой путь и медиа в личных сообщениях.
MEDIA_TEMP_DIR = media_temp_dir
MEDIA_RECOVERY_LIMIT = max(0, _env_int("STOMCHAT_MEDIA_RECOVERY_LIMIT", 5))
# Как часто доливать неразобранные снимки в очередь. Догон был однократным — при
# 745 накопленных и пяти за запуск это порядка 372 суток (14 перезапусков за 35
# суток по журналу), то есть практически никогда. Пятнадцать минут выбраны по
# пропускной способности: воркер один, разбор снимка десятки секунд, за такт
# уходит не больше MEDIA_RECOVERY_LIMIT — накопленное сходится за сутки-двое, а
# не за год, и очередь при этом не забивается.
MEDIA_RECOVERY_INTERVAL_SECONDS = max(60, _env_int("STOMCHAT_MEDIA_RECOVERY_INTERVAL", 900))
# Отметка «файл до Vision не дошёл». Пустое media_description означает «ещё не
# разбирали», и строка с ним возвращается из get_pending_media_message_ids на
# каждом запуске. Любое непустое значение снимает её с догона; отдельный текст
# (а не общее "-") нужен, чтобы отказ скачивания отличался в базе от отказа
# самого разбора.
MEDIA_UNAVAILABLE_MARK = "[медиа — файл не получен]"
_media_queue = None
_media_worker_tasks = []
_pending_albums = {}
# Что уже поставлено в очередь разбора В ЭТОМ процессе.
#
# Догон при старте (recover_pending_media_analysis) ищет в базе медиа с пустым
# media_description, а описание пишется только ПОСЛЕ разбора. Поэтому пять
# свежих снимков, которые sync_history поставил в очередь секунду назад, для
# догона выглядят необработанными — и он ставил ТЕ ЖЕ САМЫЕ ещё раз. Разбор
# идёт через платный Vision: на каждом рестарте до MEDIA_RECOVERY_LIMIT (5)
# снимков оплачивались дважды.
#
# Ограничение по длине как у HANDLED_CALLBACK_IDS: без него набор растёт на
# каждое медиа за всё время жизни процесса.
_QUEUED_MEDIA_IDS = set()
_QUEUED_MEDIA_ORDER = deque(maxlen=512)
def _remember_queued_media(msg_ids):
for queued_id in msg_ids:
if queued_id is None or queued_id in _QUEUED_MEDIA_IDS:
continue
if len(_QUEUED_MEDIA_ORDER) == _QUEUED_MEDIA_ORDER.maxlen:
_QUEUED_MEDIA_IDS.discard(_QUEUED_MEDIA_ORDER[0])
_QUEUED_MEDIA_ORDER.append(queued_id)
_QUEUED_MEDIA_IDS.add(queued_id)
# Голосовые, которые в ЭТОМ процессе уже расшифрованы или уже отданы в
# расшифровку. Отдельно от PROCESSED_MSG_IDS: тот считает обработанные АПДЕЙТЫ, а
# здесь важно другое — одно голосовое не должно уехать в Whisper дважды. Теперь
# на него смотрят два пути (живой обработчик и догон офлайн-окна), и без общей
# отметки повторная доставка апдейта или второй проход синхронизации означали бы
# второй платный вызов И вторую «🎤 Транскрипцию» в чате.
#
# Ограничение по длине как у _QUEUED_MEDIA_IDS: без него набор растёт на каждое
# голосовое за всё время жизни процесса.
_TRANSCRIBED_MSG_IDS = set()
_TRANSCRIBED_ORDER = deque(maxlen=512)
def _remember_transcribed(msg_id):
if msg_id is None or msg_id in _TRANSCRIBED_MSG_IDS:
return
if len(_TRANSCRIBED_ORDER) == _TRANSCRIBED_ORDER.maxlen:
_TRANSCRIBED_MSG_IDS.discard(_TRANSCRIBED_ORDER[0])
_TRANSCRIBED_ORDER.append(msg_id)
_TRANSCRIBED_MSG_IDS.add(msg_id)
async def get_my_id():
global MY_ID
me = await client.get_me()
MY_ID = me.id
# Меняем на числовой ID только если в конфиге реально написано 'me'
if str(config.REPORT_CHAT_ID).lower() == 'me':
config.REPORT_CHAT_ID = MY_ID
logger.info(f"✅ Отчеты будут слаться в личку (ID: {MY_ID})")
else:
# Если там число (ID группы), преобразуем в int для надежности
config.REPORT_CHAT_ID = int(config.REPORT_CHAT_ID)
logger.info(f"✅ Отчеты будут слаться в группу: {config.REPORT_CHAT_ID}")
last_summary_time = datetime.now()
def parse_state_date(value):
if not value:
return None
try:
return datetime.strptime(value, "%Y-%m-%d").date()
except ValueError:
return None
SCHEDULER_STATE_BAK_PATH = SCHEDULER_STATE_PATH + ".bak"
# Сколько дней хранить отметки о доставке отчётов. Корзины накапливались с
# первого запуска и не чистились никогда: на момент правки их 22 с 22 мая.
# Роста немного, но файл решает, отправлять ли дайджест, и разрастаться ему
# незачем.
SCHEDULER_DELIVERY_RETENTION_DAYS = 30
def _read_scheduler_file(path):
try:
with open(path, "r", encoding="utf-8") as state_file:
data = json.load(state_file)
return data if isinstance(data, dict) else None
except (OSError, ValueError):
# ValueError, а не только JSONDecodeError: повреждённый файл ловится и
# раньше разбора JSON — UnicodeDecodeError летит из самого чтения, если
# в файле оказались невалидные байты. Прежний перехват его пропускал,
# и вместо отката на резервную копию падал весь цикл планировщика.
return None
def load_scheduler_state_raw():
# Этот файл — единственное, что помнит, ушёл ли сегодняшний дайджест.
# Пустой или обрезанный файл читается как «ничего не отправляли», и отчёт
# уходит в чат ВТОРОЙ раз, поэтому при неудаче пробуем резервную копию.
state = _read_scheduler_file(SCHEDULER_STATE_PATH)
if state is None:
state = _read_scheduler_file(SCHEDULER_STATE_BAK_PATH)
if state is not None:
logger.warning("scheduler state unreadable, recovered from %s", SCHEDULER_STATE_BAK_PATH)
return state or {}
def _prune_deliveries(deliveries):
"""Отметки доставки старше SCHEDULER_DELIVERY_RETENTION_DAYS не нужны."""
if not isinstance(deliveries, dict):
return {}
cutoff = datetime.now().date() - timedelta(days=SCHEDULER_DELIVERY_RETENTION_DAYS)
kept = {}
for bucket_name, value in deliveries.items():
bucket_date = parse_state_date(bucket_name.split(":", 1)[-1])
if bucket_date is None or bucket_date >= cutoff:
kept[bucket_name] = value
return kept
def load_scheduler_state():
state = load_scheduler_state_raw()
if not state:
return None, None
return (
parse_state_date(state.get("last_daily_date")),
parse_state_date(state.get("last_weekly_date")),
)
def save_scheduler_state(last_daily_date, last_weekly_date, deliveries=None):
if deliveries is None:
deliveries = load_scheduler_state_raw().get("deliveries", {})
state = {
"last_daily_date": last_daily_date.isoformat() if last_daily_date else None,
"last_weekly_date": last_weekly_date.isoformat() if last_weekly_date else None,
"deliveries": _prune_deliveries(deliveries),
}
temp_path = SCHEDULER_STATE_PATH + ".tmp"
try:
# os.replace атомарна, но без fsync содержимое временного файла может не
# дойти до диска раньше переименования: после сбоя питания на месте
# состояния оказывается пустой файл, «дайджест не отправлялся», и отчёт
# уходит в чат второй раз. Та же связка, что уже стоит в
# assistant.save_state: временный файл, fsync, замена, резервная копия.
with open(temp_path, "w", encoding="utf-8") as state_file:
json.dump(state, state_file, ensure_ascii=False, indent=2)
state_file.flush()
os.fsync(state_file.fileno())
if os.path.exists(SCHEDULER_STATE_PATH):
try:
os.replace(SCHEDULER_STATE_PATH, SCHEDULER_STATE_BAK_PATH)
except OSError as backup_err:
logger.warning("scheduler state backup failed: %s", backup_err)
os.replace(temp_path, SCHEDULER_STATE_PATH)
return True
except OSError as write_err:
# Молча провалить запись нельзя: в памяти день уже помечен отправленным,
# а после перезапуска отчёт уйдёт повторно.
logger.error("SCHEDULER STATE NOT SAVED: %s — отчёт может уйти повторно", write_err)
try:
if os.path.exists(temp_path):
os.remove(temp_path)
except OSError:
pass
return False
def target_delivery_key(chat_id, topic_id):
topic = "main" if topic_id is None else str(topic_id)
return f"{chat_id}:{topic}"
def delivery_bucket(report_kind, report_date):
return f"{report_kind}:{report_date.isoformat()}"
def load_sent_targets(report_kind, report_date):
deliveries = load_scheduler_state_raw().get("deliveries", {})
bucket = deliveries.get(delivery_bucket(report_kind, report_date), {})
if not isinstance(bucket, dict):
return set()
return {target_key for target_key, value in bucket.items() if value}
def mark_target_delivered(report_kind, report_date, target_key, last_daily_date, last_weekly_date, message_id=None):
state = load_scheduler_state_raw()
deliveries = state.get("deliveries", {})
if not isinstance(deliveries, dict):
deliveries = {}
bucket_name = delivery_bucket(report_kind, report_date)
bucket = deliveries.setdefault(bucket_name, {})
bucket[target_key] = {
"delivered_at": datetime.now().isoformat(timespec="seconds"),
"message_id": message_id,
}
save_scheduler_state(last_daily_date, last_weekly_date, deliveries)
def parse_status_utc(value):
if not value:
return None
try:
return datetime.fromisoformat(value).astimezone(timezone.utc)
except ValueError:
return None
async def runtime_telemetry_task():
while True:
try:
runtime_guard.write_heartbeat("runtime_telemetry")
if psutil is None:
logger.info("runtime_memory psutil_unavailable")
else:
process = psutil.Process(os.getpid())
info = process.memory_info()
try:
full_info = process.memory_full_info()
except Exception:
full_info = info
private_bytes = (
getattr(full_info, "private", None)
or getattr(full_info, "uss", None)
or getattr(info, "rss", 0)
)
try:
open_files = len(process.open_files())
except Exception:
open_files = -1
logger.info(
"runtime_memory pid=%s rss_mb=%.2f private_mb=%.2f vms_mb=%.2f threads=%s open_files=%s",
os.getpid(),
getattr(info, "rss", 0) / 1024 / 1024,
private_bytes / 1024 / 1024,
getattr(info, "vms", 0) / 1024 / 1024,
process.num_threads(),
open_files,
)
except Exception as exc:
logger.warning("runtime_memory_error %s", exc)
cleanup_temp_media()
await asyncio.sleep(900)
TEMP_MEDIA_MAX_AGE_SECONDS = 6 * 3600
def cleanup_temp_media(max_age_seconds=TEMP_MEDIA_MAX_AGE_SECONDS):
"""
Подметает temp_media от файлов, переживших свою обработку.
Штатные пути уборки есть, но они не покрывают обрыв download_media по
таймауту (файл уже создан, а путь наверх не вернулся), падение извлечения
кадра и убийство процесса сторожем. Чистки по расписанию не было вообще:
на момент добавления в каталоге лежало 69 файлов на 43.6 МБ, включая
13 нулевых, самый старый — почти полугодовой давности.
Уборка по возрасту, а не по имени: так покрываются все пути утечки сразу,
и активная обработка не задевается — 6 часов сильно больше любого таймаута.
"""
removed = 0
freed = 0
try:
# Каталог берём из MEDIA_TEMP_DIR, а не из литерала. Переменная
# STOMCHAT_MEDIA_TEMP_DIR соблюдалась ТОЛЬКО на скачивании: уборщик
# подметал пустой "temp_media", а файлы копились в настроенном каталоге
# вечно — вместе с обрывками скачиваний и голосовыми.
if not os.path.isdir(MEDIA_TEMP_DIR):
return 0
cutoff = time.time() - max_age_seconds
for name in os.listdir(MEDIA_TEMP_DIR):
path = os.path.join(MEDIA_TEMP_DIR, name)
try:
if not os.path.isfile(path) or os.path.getmtime(path) > cutoff:
continue
size = os.path.getsize(path)
os.remove(path)
removed += 1
freed += size
except OSError:
# Файл может быть занят активной обработкой — заберём в следующий раз.
continue
except Exception as exc:
logger.warning("temp_media_cleanup_error %s", exc)
if removed:
logger.info(f"temp_media cleanup: removed {removed} stale files, freed {freed/1e6:.1f} MB")
return removed
async def heartbeat_task():
# Единственный фоновый цикл, у которого не было try/except (остальные
# обёрнуты). write_heartbeat делает os.replace и на Windows ловит
# PermissionError, если файл в этот момент держит антивирус/индексатор;
# после 5 ретраев он поднимает OSError. Таск умирал навсегда, heartbeat
# переставал обновляться, и через WATCHDOG_STALE_SECONDS сторож убивал
# процесс — как правило, посреди генерации саммари или анализа снимка.
while True:
try:
runtime_guard.write_heartbeat("heartbeat")
except asyncio.CancelledError:
raise
except Exception as e:
logger.warning(f"Heartbeat write failed, continuing: {e}")
await asyncio.sleep(runtime_guard.HEARTBEAT_INTERVAL_SECONDS)
async def summary_watchdog_task():
while True:
await asyncio.sleep(SUMMARY_STATUS_CHECK_SECONDS)
try:
status = runtime_guard.read_summary_status()
if not status.get("active"):
continue
updated_at = parse_status_utc(status.get("utc"))
if not updated_at:
continue
age = (datetime.now(timezone.utc) - updated_at).total_seconds()
if age <= SUMMARY_STALE_SECONDS:
continue
logger.error(
"summary watchdog forcing restart: stage=%s kind=%s chat=%s age=%.1fs status=%s",
status.get("stage"),
status.get("kind"),
status.get("chat_id"),
age,
status,
)
runtime_guard.dump_runtime_state("summary_watchdog_stale")
os._exit(79)
except Exception:
logger.exception("summary watchdog failed")
# Час, с которого начинается дневное окно. Он НЕ независим от
# config.REPORT_HOUR: пока это были две несвязанные константы, стык окон
# держался на совпадении чисел. При REPORT_HOUR=22 (текущий бой) окна
# перекрываются на 2 часа — по локальному снимку это 15.8% трафика, то есть
# вечер уходит в дайджест ДВАЖДЫ, но не теряется. А при REPORT_HOUR меньше 20
# (в config.example.py стоит 0) между выпусками возникает дыра: окно кончается в
# REPORT_HOUR, а следующее начинается только в 20:00 того же дня — до 20 часов
# переписки в сутки не попадают ни в один выпуск. daily_window_start() связывает
# их через last_sent_date, поэтому дыра закрывается при любом REPORT_HOUR.
DIGEST_WINDOW_START_HOUR = 20
# Насколько глубоко догоняем пропущенные выпуски. Без ограничения месяц простоя
# уехал бы одним запросом в промпт: по локальному снимку это порядка 8000
# сообщений против 228 в обычные сутки (худшая измеренная тройка дней — 1458).
# Ограничение делает потерю видимой в журнале вместо того, чтобы отдать
# генерации заведомо неподъёмный лог; за пределами трёх суток период всё равно
# покрывает недельный отчёт.
DIGEST_CATCHUP_MAX_DAYS = 3
WEEKLY_REPORT_HOUR = 10
WEEKLY_WINDOW_DAYS = 7
# Догон недельного отчёта тоже ограничен: 10 суток вместо 7 — это запас на
# пропущенный понедельник, а не бесконечное окно.
WEEKLY_CATCHUP_MAX_DAYS = 10
# Двух выпусков за трое суток быть не должно. Без этой границы догон, отработав
# в воскресенье, выпускал бы вторую газету в понедельник — на почти том же
# материале и с отдельной страницей Telegraph.
WEEKLY_MIN_GAP_DAYS = 3
def daily_window_start(now, last_sent_date):
"""
Нижняя граница окна дневного дайджеста.
Раньше она считалась только от now: «(now - 1 день) в 20:00». Пока выпуски
идут каждый день, этого хватает, но last_sent_date хранится ровно для
другого случая — пропущенного дня. Если вчерашний дайджест не ушёл (бот
лежал, цель отвалилась, сообщений было меньше порога), то сегодняшнее окно
всё равно начиналось с 20:00 вчера, и промежуток от конца последнего выпуска
(позавчера, REPORT_HOUR) до вчерашних 20:00 — порядка 22 часов — не попадал
НИ В ОДИН выпуск. Добор по is_summarized этого не спасает: он включается
только когда в окне меньше min_count сообщений, и берёт лишь недостачу.
Границы остаются наивными локальными: в UTC их переводит database._date_text,
и подменять эту трактовку здесь нельзя.
"""
base_start = (now - timedelta(days=1)).replace(
hour=DIGEST_WINDOW_START_HOUR, minute=0, second=0, microsecond=0
)
if not last_sent_date:
return base_start
# Прошлый выпуск закончился не раньше REPORT_HOUR своего дня — оттуда и
# продолжаем. Небольшое перекрытие лучше пропуска: повтор врачи увидят,
# потерянные сутки — нет.
resume_from = datetime.combine(last_sent_date, datetime.min.time()).replace(
hour=int(config.REPORT_HOUR) % 24
)
floor_start = now - timedelta(days=DIGEST_CATCHUP_MAX_DAYS)
return max(floor_start, min(base_start, resume_from))
def weekly_report_due(now, last_weekly_date):
"""
Пора ли выпускать недельный отчёт.
Условие было только «понедельник, 10:00, сегодня ещё не отправляли».
Пропущенный понедельник — упавший бот, недоставленная цель — терял выпуск
НАВСЕГДА: следующий понедельник берёт окно в 7 дней и до пропущенной недели
не достаёт. last_weekly_date хранится в том же файле состояния, что и
дневная дата, и здесь он и нужен.
"""
if now.hour < WEEKLY_REPORT_HOUR:
return False
if last_weekly_date == now.date():
return False
if last_weekly_date is not None and (now.date() - last_weekly_date).days < WEEKLY_MIN_GAP_DAYS:
return False
if now.weekday() == 0:
return True
# Догон: неделя с последнего выпуска прошла, а понедельника мы не увидели.
if last_weekly_date is None:
return False
return (now.date() - last_weekly_date).days >= WEEKLY_WINDOW_DAYS
def weekly_window_start(now, last_weekly_date):
"""Нижняя граница недельного окна — с догоном по last_weekly_date."""
base_start = now - timedelta(days=WEEKLY_WINDOW_DAYS)
if not last_weekly_date:
return base_start
resume_from = datetime.combine(last_weekly_date, datetime.min.time()).replace(
hour=WEEKLY_REPORT_HOUR
)
floor_start = now - timedelta(days=WEEKLY_CATCHUP_MAX_DAYS)
return max(floor_start, min(base_start, resume_from))
def resolve_report_targets():
"""Цели рассылки из конфига — с громким отказом вместо молчаливого нуля.
Прежде здесь стоял голый `except: targets = []`. Любая ошибка — и планировщик
поднимался с нулём целей: в журнале бодрое «Планировщик активен. Целей: 0»,
а ни один врач не получал ни дайджеста, ни недельной сводки, и узнать об этом
было неоткуда. Ошибку формы конфиг вообще не ловит: `json.loads` на
`REPORT_TARGETS={"chat_id": -100}` отдаёт валидный dict, а не список.
"""
# Формат в .env: REPORT_TARGETS=[{"chat_id": -100123, "topic_id": null}, {"chat_id": -100456, "topic_id": 390}]
try:
raw = config.REPORT_TARGETS
except Exception as exc:
logger.error(
"REPORT_TARGETS не прочитан (%s: %s) — ни дайджест, ни недельная "
"сводка не уйдут НИКОМУ",
type(exc).__name__, exc, exc_info=True,
)
return []
if not isinstance(raw, list):
logger.error(
"REPORT_TARGETS не список, а %s — рассылка не уйдёт НИКОМУ. "
'Ожидается [{"chat_id": -100123, "topic_id": null}]',
type(raw).__name__,
)
return []
targets = []
for position, target in enumerate(raw):
# Битая цель уносила ВСЮ рассылку: в `.test` обработчике элемент-строка
# падал на `target.get` с AttributeError. Пропускаем ровно её, остальные
# чаты сводку получают.
if not isinstance(target, dict) or "chat_id" not in target:
logger.error(
"REPORT_TARGETS[%d] пропущен (%r): нет chat_id — этот чат "
"сводок не получит",
position, target,
)
continue
targets.append(target)
if raw and not targets:
logger.error(
"REPORT_TARGETS: все %d целей битые — рассылка не уйдёт НИКОМУ",
len(raw),
)
return targets
async def scheduler_task(bot_client):
"""Рассылка по всем целям из конфига."""
targets = resolve_report_targets()
logger.info(f"📅 Планировщик активен. Целей: {len(targets)}")
last_sent_date, last_weekly_date = load_scheduler_state()
# Кэш сгенерированного дайджеста живёт МЕЖДУ кругами цикла, а не внутри
# одного: см. комментарий на присваивании generated_cache ниже.
daily_cache_date = None
daily_cache_text = None
while True:
try:
now = datetime.now()
# 1. ЕЖЕДНЕВНЫЙ ДАЙДЖЕСТ (Daily)
# Проверка времени (REPORT_HOUR) и того, что сегодня еще не отправляли
if now.hour >= config.REPORT_HOUR and last_sent_date != now.date():
# Окно: от конца прошлого выпуска (либо 20:00 вчера) до сейчас.
end_time = now
start_time = daily_window_start(now, last_sent_date)
logger.info(
"Daily окно: %s -> %s (last_sent=%s)",
start_time, end_time, last_sent_date,
)
messages = await asyncio.wait_for(
database.get_messages_for_daily_summary(start_time, end_time, min_count=100),
timeout=30,
)
if messages:
logger.info(f"🔥 Daily контент готов ({len(messages)} шт). Рассылка...")
# Кэш для текста (чтобы генерировать 1 раз на все чаты).
#
# Он обязан переживать круг цикла. Пока generated_cache
# создавался здесь заново, сломанная цель означала ПОЛНУЮ
# повторную генерацию каждые 10 минут: last_sent_date не
# продвигается, пока не доставлено во все цели, поэтому
# следующий заход снова шёл в генерацию — платный вызов LLM
# и новая страница Telegraph на каждый круг, до полуночи.
if daily_cache_date != now.date():
daily_cache_date = now.date()
daily_cache_text = None
generated_cache = daily_cache_text
sent_targets = load_sent_targets("daily", now.date())
target_keys = [
target_delivery_key(target.get('chat_id'), target.get('topic_id'))
for target in targets
if target.get('chat_id')
]
# Проходим по всем целям
for target in targets:
tgt_chat = target.get('chat_id')
tgt_topic = target.get('topic_id')
if not tgt_chat: continue
tgt_key = target_delivery_key(tgt_chat, tgt_topic)
if tgt_key in sent_targets:
logger.info("Daily target already delivered; skip duplicate target=%s", tgt_key)
continue
try:
logger.info(f"📤 Отправка Daily в {tgt_chat} (Topic: {tgt_topic})...")
async def daily_delivery_hook(sent_message, target_key=tgt_key):
sent_targets.add(target_key)
mark_target_delivered(
"daily",
now.date(),
target_key,
last_sent_date,
last_weekly_date,
getattr(sent_message, "id", None),
)
# Передаем кэш и сохраняем результат
result_text = await summarizer.process_summary_batch(
messages,
bot_client,
chat_id=tgt_chat,
topic_id=tgt_topic,
msg_count=len(messages),
cached_message=generated_cache,
delivery_hook=daily_delivery_hook,
)
# Если генерация прошла успешно, запоминаем текст для следующих кругов
if result_text:
if tgt_key not in sent_targets:
sent_targets.add(tgt_key)
mark_target_delivered("daily", now.date(), tgt_key, last_sent_date, last_weekly_date)
if not generated_cache:
generated_cache = result_text
daily_cache_text = result_text
except Exception:
logger.exception(f"Daily send failed chat={tgt_chat}")
if target_keys and all(target_key in sent_targets for target_key in target_keys):
# Помечаем сообщения прочитанными 1 раз после всех рассылок
msg_ids = [m[0] for m in messages]
await asyncio.wait_for(database.mark_messages_as_summarized(msg_ids), timeout=30)
last_sent_date = now.date()
save_scheduler_state(last_sent_date, last_weekly_date)
logger.info("✅ Ежедневная рассылка завершена.")
else:
missing_targets = [target_key for target_key in target_keys if target_key not in sent_targets]
logger.error("Daily was not delivered to all targets; missing=%s messages remain unsummarized.", missing_targets)
# 2. ЕЖЕНЕДЕЛЬНАЯ ГАЗЕТА (Weekly)
# Запуск: Понедельник (weekday == 0), 10:00 утра — либо догон, если
# понедельник был пропущен (см. weekly_report_due).
if weekly_report_due(now, last_weekly_date):
logger.info(
"🗞 Наступило время Weekly отчета (weekday=%s, last_weekly=%s)...",
now.weekday(), last_weekly_date,
)
# Период: последние 7 полных дней, с догоном по last_weekly_date
end_weekly = now
start_weekly = weekly_window_start(now, last_weekly_date)
# Получаем сообщения за диапазон
weekly_messages = await asyncio.wait_for(
database.get_messages_for_range(start_weekly, end_weekly),
timeout=30,
)
if weekly_messages:
logger.info(f"💎 Weekly контент готов ({len(weekly_messages)} шт). Рассылка...")
weekly_sent_targets = load_sent_targets("weekly", now.date())
weekly_target_keys = [
target_delivery_key(target.get('chat_id'), target.get('topic_id'))
for target in targets
if target.get('chat_id')
]
for target in targets:
tgt_chat = target.get('chat_id')
tgt_topic = target.get('topic_id')
if not tgt_chat: continue
tgt_key = target_delivery_key(tgt_chat, tgt_topic)
if tgt_key in weekly_sent_targets:
logger.info("Weekly target already delivered; skip duplicate target=%s", tgt_key)
continue
try:
logger.info(f"📤 Отправка Weekly в {tgt_chat} (Topic: {tgt_topic})...")
async def weekly_delivery_hook(sent_message, target_key=tgt_key):
weekly_sent_targets.add(target_key)
mark_target_delivered(
"weekly",
now.date(),
target_key,
last_sent_date,
last_weekly_date,
getattr(sent_message, "id", None),
)
result_text = await summarizer.process_weekly_batch(
weekly_messages,
bot_client,
chat_id=tgt_chat,
topic_id=tgt_topic,
delivery_hook=weekly_delivery_hook,
)
if result_text:
if tgt_key not in weekly_sent_targets:
weekly_sent_targets.add(tgt_key)
mark_target_delivered("weekly", now.date(), tgt_key, last_sent_date, last_weekly_date)
except Exception:
logger.exception(f"Weekly send failed chat={tgt_chat}")
if weekly_target_keys and all(target_key in weekly_sent_targets for target_key in weekly_target_keys):
last_weekly_date = now.date()
save_scheduler_state(last_sent_date, last_weekly_date)
logger.info("✅ Еженедельная рассылка (Weekly) завершена.")
else:
missing_targets = [target_key for target_key in weekly_target_keys if target_key not in weekly_sent_targets]
logger.error("Weekly was not delivered to all targets; missing=%s scheduler state not advanced.", missing_targets)
await asyncio.sleep(600) # Проверка каждые 10 минут
except Exception as e:
logger.error(f"Ошибка планировщика: {e}")
await asyncio.sleep(60)
# Потолок на один проход пингов. Своих таймаутов у этих двух вызовов не было ни
# одного: зависший send_message (Telethon держит request_retries=10 и спит на
# FloodWait) останавливал ВЕСЬ цикл пингов навсегда и совершенно молча — таск
# жив, sleep(3600) не наступает, в журнале ни строки.
#
# Значение выведено из бюджета законной работы, а не выбрано отдельно:
# MAX_PINGS_PER_CYCLE = 5 пингов, на каждый до 60 с генерации (timeout=60 в
# assistant) плюс отправка с ретраями и сном на FloodWait — порядка 600 с.
# Двойной запас даёт 1200 с. Даже если оба прохода выберут потолок целиком,
# 2400 с меньше часового сна, то есть почасовая частота пингов сохраняется.
PING_PHASE_TIMEOUT_SECONDS = 1200
async def pm_ping_scheduler_task(bot_client):
"""Задача периодической проверки неактивности пользователей в ЛС и отправки им пинга."""
logger.info("📅 Планировщик пингов в ЛС активен.")
while True:
try:
await asyncio.wait_for(
assistant.check_and_send_pm_pings(bot_client),
timeout=PING_PHASE_TIMEOUT_SECONDS,
)
except asyncio.TimeoutError:
logger.error(
"PM pings timed out after %ss — проход прерван, цикл продолжается",
PING_PHASE_TIMEOUT_SECONDS,
)
except Exception as e:
logger.error(f"Error in pm_ping_scheduler_task (PM pings): {e}")
await asyncio.sleep(3600) # Проверка каждый час
# Порог, после которого telethon перестаёт спать на FloodWait и поднимает
# исключение. По умолчанию он равен 60 с и НЕ задавался, а при request_retries=10
# это до 600 секунд сна ВНУТРИ одного await send_message — молча, потому что
# строка про сон идёт уровнем INFO у логгера telethon.client.users, а
# runtime_guard приглушает telethon до ERROR. Десять минут внутри одного вызова
# означают десять минут под удерживаемым замком диалога: следующее сообщение
# того же врача встаёт в очередь за ним.
#
# СВЕРЕНО С ИСХОДНИКАМИ БИБЛИОТЕКИ, не по памяти (telethon 1.42.0):
# client/telegrambaseclient.py:254 flood_sleep_threshold: int = 60 — да, 60;
# client/telegrambaseclient.py:249 request_retries: int = 5 — библиотечный
# дефолт пять, десять ставим МЫ ниже при создании клиентов;
# client/users.py:70 for attempt in retry_range(self._request_retries)
# client/users.py:120 if e.seconds <= self.flood_sleep_threshold:
# client/users.py:122 await asyncio.sleep(e.seconds)
# То есть сон действительно происходит ВНУТРИ одной попытки внутри одного await,
# логируется через self._log[__name__] на INFO, и повторяется до request_retries
# раз. Прежде это было записано как «верю агенту, помечаю непроверенным» —
# теперь проверено по коду. Оговорка сохраняется: «10 x (30 + 20) = 500 с» это
# ВЕРХНЯЯ ОЦЕНКА, потому что 30 с приходят из транспортного слоя, а не из этого
# цикла; сам цикл гарантирует только «до request_retries снов по порогу».
#
# Ноль здесь был бы хуже, а не лучше: FloodWait стал бы прилетать во все ~120
# вызовов assistant.py, где стоит голый except, и короткая задержка в пять секунд
# из «медленно, но доставлено» превратилась бы в «молча не доставлено». Поэтому
# порог небольшой, но не нулевой: короткое ожидание пересиживаем, длинное отдаём
# вызывающему. Полную границу по времени даёт только бюджет на вызове —
# tg_safety.guard; пока он подключён не везде, это ограничение сверху, а не
# гарантия.
#
# ДОЛГ P2 врач: в main.py tg_safety не подключён НИ В ОДНОЙ точке -> замер
# 29 июля 2026 разбором ast: 11 вызовов Telegram, из них 0 через tg_safety и 5
# меняющих чат (:1352 публикация расшифровки голосового, :1437 отказ
# распознавания, :2010 подтверждение закладки, :2082 и :2086 удаление
# сообщений); у каждого голого await верхняя оценка сна внутри telethon
# 10 x (30 + 20) = 500 с при живом клиенте, и врач видит зависший бот без
# единой строки в журнале. Для сравнения: assistant.py прошёл на tg_safety 11
# вызовов из 93 (12%), summarizer.py — 0 из 4. Закрывать не здесь: перевод
# отправок на tg_safety — работа лана доставки, а этот файл только объявляет
# порог, поэтому долг помечен, а не починен.
TELETHON_FLOOD_SLEEP_THRESHOLD = max(0, _env_int("STOMCHAT_FLOOD_SLEEP_THRESHOLD", 20))
# 1. Клиент Юзербота (Твой аккаунт) - только слушает
client = TelegramClient(
config.SESSION_NAME,
config.API_ID,
config.API_HASH,
timeout=30,
request_retries=10,
connection_retries=1000,
retry_delay=5,
auto_reconnect=True,
flood_sleep_threshold=TELETHON_FLOOD_SLEEP_THRESHOLD,
)
# 2. Клиент Бота - только пишет и крепит
bot_client = TelegramClient(
'bot_session',
config.API_ID,
config.API_HASH,
timeout=30,
request_retries=10,
connection_retries=1000,
retry_delay=5,
auto_reconnect=True,
flood_sleep_threshold=TELETHON_FLOOD_SLEEP_THRESHOLD,
)
# Wrapper to track bot's own outgoing message IDs for safety wipe commands
original_send_message = bot_client.send_message
async def patched_send_message(*args, **kwargs):
text = None
if len(args) > 1:
text = args[1]
elif "message" in kwargs:
text = kwargs.get("message")
has_file = bool(kwargs.get("file") or (len(args) > 2 and args[2]))
if not has_file:
if not text or not (str(text) if text is not None else "").strip():
entity = args[0] if len(args) > 0 else kwargs.get("entity")
logger.warning(f"bot_client.send_message: aborted sending empty message to entity={entity}")
return None
sent_msg = await original_send_message(*args, **kwargs)
if sent_msg and hasattr(sent_msg, 'id') and hasattr(sent_msg, 'peer_id'):
try:
# Пересчёт peer -> chat_id отдаём самой библиотеке. Ручная
# арифметика здесь была верной (проверено на PeerChannel/PeerChat/
# PeerUser), но она повторяет utils.get_peer_id и молча разойдётся
# с ним, если Telethon добавит новый тип peer.
chat_id = telethon_utils.get_peer_id(sent_msg.peer_id)
if chat_id:
await database.save_bot_sent_message(sent_msg.id, chat_id)
except Exception as e:
logger.error(f"Error saving bot outgoing message ID: {e}")
return sent_msg
bot_client.send_message = patched_send_message
def start_media_analysis_workers():
global _media_queue, _media_worker_tasks
if _media_queue is None:
_media_queue = asyncio.Queue(maxsize=MEDIA_QUEUE_MAX_SIZE)
_media_worker_tasks = [task for task in _media_worker_tasks if not task.done()]
while len(_media_worker_tasks) < MEDIA_WORKER_COUNT:
worker_id = len(_media_worker_tasks) + 1
_media_worker_tasks.append(
runtime_guard.create_task(media_analysis_worker(worker_id), f"media_analysis_{worker_id}")
)
async def stop_media_analysis_workers():
global _media_worker_tasks
if not _media_worker_tasks:
return
for task in _media_worker_tasks:
task.cancel()
await asyncio.gather(*_media_worker_tasks, return_exceptions=True)
_media_worker_tasks = []
async def enqueue_media_analysis(messages, msg_id, text, media_type_hint=None, bulk=False):
"""
Ставит медиа в очередь разбора. Возвращает True, если место нашлось.
bulk=True — постановка пачкой из догоняющей синхронизации: там переполнение
очереди штатно и не должно писать строку ERROR на каждый снимок.
Непоставленные никуда не пропадают: в базе у них пустое media_description,
и их подбирает recover_pending_media_analysis при следующих запусках.
"""
# Вызываем безусловно, а не только при _media_queue is None.
# start_media_analysis_workers() идемпотентна: она отбрасывает завершившиеся
# таски и добирает недостающие. Раньше её звали лишь на старте, поэтому
# умерший воркер (по умолчанию он один) не поднимался никогда: очередь молча
# заполнялась до предела, и дальше ВСЁ медиа уходило в logger.error без
# анализа — до следующего перезапуска процесса.
start_media_analysis_workers()
try:
_media_queue.put_nowait((messages, msg_id, text, media_type_hint))
# Запоминаем ВСЕ сообщения пачки, а не только msg_id: у альбома в базе
# своя строка на каждый снимок, и догон нашёл бы остальные по пустому
# описанию. Отмечаем только после успешной постановки — непоставленные