Reply Pilot Worker

Tato stranka popisuje Java 21 modul reply-pilot-worker/.

Role modulu

  • background job container mimo request cyklus
  • owner mailbox import workflow nad reply-pilot-be
  • umi periodicky obnovovat Gmail watch a volitelne cist Gmail Pub/Sub Pull subscription nad reply-pilot-be
  • umi drzet sjednoceny gmail_synchronizer job, ktery spojuje watch pull a mailbox import do jedne pipeline
  • umi kazdych 120 sekund spoustet uplnou CME company synchronizaci nad reply-pilot-be
  • kazde dve hodiny autentizovane triggeruje logical dump v reply-pilot-db-backup; nezna DB host, user ani password
  • umi kazdou minutu synchronizovat lokalni task_jira cache podle Jiry pres reply-pilot-be
  • umi volitelne periodicky spoustet AI klasifikaci prefiltrovanych requirement kandidatu nad reply-pilot-be
  • umi volitelne periodicky spoustet finalni requirement agregaci nad reply-pilot-be
  • pripravena hranice pro dalsi async joby a retry flow

Runtime

  • worker nema verejne HTTP rozhrani na hostu
  • worker vystavuje interni HTTP status API na sdilene Docker siti:
  • GET /healthz pro healthcheck kontejneru
  • GET /statusz pro backend monitoring
  • lokalni heartbeat zustava v data/worker-heartbeat.json; GET /healthz jeho cerstvost vyuziva pro self-check workeru
  • heartbeat ted krom last_job drzi i job_metrics po jednotlivych workerech a pole running_jobs, aby slo delat jednoduchy ops report nad AI a evaluation pipeline
  • loguje vyhradne na stdout/stderr; lokalne pouzij docker compose logs
  • backend vola pres BACKEND_API_BASE_URL, defaultne http://reply-pilot-be:5000
  • worker pouziva Spring ThreadPoolTaskScheduler v jednom procesu:
  • sedm samostatnych @Component jobu ma @Scheduled(fixedDelay); interval se pocita od dokonceni predchoziho behu, ne od jeho zacatku
  • Gmail watch a synchronizer maji Spring Trigger, ktery zachovava dynamicke retry deadlines, exponencialni backoff a intervaly mailbox pollingu
  • devet jobu ma k dispozici devet execution vlaken a heartbeat dalsi vlakno
  • jeden job se v teto instanci nespusti podruhe, dokud jeho predchozi beh stale trva
  • gmail_synchronizer je jediny owner mailbox flow a drzi si vlastni stav rozbehnuteho importu
  • produkcne bezi jedna instance scheduleru; nevytvari se distribuovany scheduler
  • HTTP volani pouzivaji JDK HttpClient, status API JDK HttpServer, JSON Jackson
  • heartbeat se zapisuje atomicky; pri selhani zapisu se jeho timestamp neobnovi
  • SIGTERM rusi cekajici behy, prerusi aktivni HTTP requesty a ukonci status server
  • Java healthcheck java -jar /app/reply-pilot-worker.jar --healthcheck pouze vola lokalni /healthz; nevyzaduje Python ani curl

Kod a overeni

Vsechny tridy jsou v src/main/java/cz/replypilot/worker/:

  • ReplyPilotWorker.java: Spring context a jednorazove spusteni
  • WorkerScheduling.java: @EnableScheduling, pool a registrace Gmail triggers
  • CmeCompanySyncJob.java: anotovany CME job volajici /api/cme/company-sync/run
  • DatabaseBackupJob, OrganizeMeetingResumeJob, WorkWizardSnoozeCleanupJob, JiraTaskSyncJob, RequirementAiJob, RequirementEvaluationJob: ostatni periodicke @Scheduled joby
  • GmailWatchJob.java, GmailSynchronizerJob.java: Gmail watch a mailbox pipeline
  • WorkerJob.java: spolecne provedeni jobu a zaznam chyb; GmailJob.java: retry trigger
  • WorkerStatus.java: heartbeat, running jobs a metriky
  • WorkerStatusServer.java: interni /healthz a /statusz
  • BackendClient.java: HTTP/JSON kontrakt, timeouty a redakce credentials v chybach
  • WorkerConfig.java: existujici env defaults a lokalni dotenv precedence
  • WorkerApplication.java: trvaly proces, --once a --healthcheck
  • z korene repozitare: mvn -f reply-pilot-worker/pom.xml test package

--once neregistruje scheduling ani status server. On-start prepinace a konfiguracni defaults zustavaji stejne; Gmail synchronizer a oba pevne maintenance joby startuji okamzite. CME je ve vychozim stavu vypnuty, po zapnuti ceka prvni interval, pokud neni WORKER_CME_SYNC_ON_START=true.

Aktualni job

  • gmail_synchronizer v kratkem intervalu vola POST /api/mailbox/watch/pull; pokud notifikace ukazuji na zmenu mailboxu, zaqueueuje import
  • jednou za WORKER_IMPORT_POLL_INTERVAL_SECONDS si gmail_synchronizer kontroluje GET /api/mailbox/import/status
  • kdyz import je QUEUED nebo RUNNING, gmail_synchronizer ho po krocich dotahuje pres POST /api/mailbox/import/step
  • auto-queue importu zustava volitelne pres WORKER_IMPORT_AUTO_QUEUE_ENABLED=true; kdyz je zapnute, gmail_synchronizer pravidelne frontuje incremental import i bez watch notifikace, ne opakovany full scan
  • pro aktivni import gmail_synchronizer opakovane vola POST /api/mailbox/import/step, dokud backend nevrati COMPLETED nebo FAILED
  • POST /api/cme/company-sync/run jednou za WORKER_CME_SYNC_INTERVAL_SECONDS; synchronizuje dodavatele z dodavatel včetně odpovědného obchodníka z dodavatel_informace a rezervace z osloveni_dodavatele; firmy páruje výhradně podle povinného platného IČO a při více aktivních kandidátech zapíše chybu bez automatického vítěze
  • POST http://reply-pilot-db-backup:8080/api/backups jednou za WORKER_DATABASE_BACKUP_INTERVAL_SECONDS (default 7200); trigger pouziva DATABASE_BACKUP_API_TOKEN, bezi i on-start a worker neceka na dokonceni dumpu
  • POST /api/jira/task-sync/run jednou za WORKER_JIRA_TASK_SYNC_INTERVAL_SECONDS; backend si pamatuje Jira updated watermark v public.jira_sync_state a lokalni task_jira bere jako cache
  • POST /api/tasks/organize-meeting/resume-due pri startu a kazdych 300 sekund
  • POST /api/work-wizard/snoozes/cleanup pri startu a kazdych 900 sekund
  • POST /api/requirements/ai-classify/run v konfigurovatelnem intervalu; default je vypnuty, aby se AI job nepoustel bez vedome konfigurace
  • POST /api/requirements/evaluate/run v konfigurovatelnem intervalu; typicky navazuje na deterministic a AI vrstvu a prepocitava finalni company atributy
  • Gmail modul pri POST /api/mailbox/watch/pull ulozi pending ack ID a potvrdi je az po uspesnem mailbox-wide snapshot/history zapisu; navazany incremental import pak ridi gmail_synchronizer
  • heartbeat se udrzuje i v idle stavu, aby worker mel lokalni self-observed stav; backend uz ho necte pres sdileny mount, ale pres GET /statusz
  • pri timeoutu nebo backend chybe konkretni job nesmi shodit cely worker proces; chyba se materializuje do job_metrics a scheduler pokracuje dalsimi joby