import atexit from dataclasses import dataclass, field from datetime import datetime, timezone from pathlib import Path from typing import Any, ClassVar, Iterator, NamedTuple import uuid import gradio as gr from injector import inject from src.backend.application.chat_completion import ( ChatCompletionAnswerEvent, ChatCompletionDataPointsEvent, ChatCompletionInput, ChatCompletionUseCase, QueryGeneratedEvent, TokenUsageEvent, ) from src.backend.config.settings import CARINA_INDEX_NAME from src.backend.dependency_injector import create_app_injector, create_injector from src.backend.domain.models import InputFile, PdfParser, Property from src.backend.domain.models.chat_completion import Content, Message, Messages, Role from src.backend.domain.models.access_log import AccessLogEntry from src.backend.domain.models.chat_log import ChatLogEntry, FeedbackCategory, FeedbackLog from src.backend.domain.models.data_structure import DataPoint from src.backend.domain.models.llm import ModelName from src.backend.domain.services.i_chat_log_repository import IChatLogRepository from src.backend.domain.services.i_file_processor import IFileProcessorService from src.backend.domain.services.i_user_repository import IUserRepository from src.backend.domain.services.tool import Tool from src.backend.utils.ip_allowlist import ( extract_client_ip, is_ip_allowed, load_allowed_ips, ) ASSETS_DIR = Path(__file__).resolve().parent / "assets" MODEL_CHOICES: dict[str, ModelName] = { "gemini-3-flash-preview(軽量)": ModelName.GEMINI_3_FLASH_PREVIEW, } DEFAULT_MODEL_LABEL: str = next(iter(MODEL_CHOICES)) class ProgressPanel(NamedTuple): # チャット上で折りたたみ表示する途中経過(検索)の1パネル。 # panel_id / parent_id を指定すると Gradio の metadata でトグルを入れ子にできる。 title: str content: str panel_id: str | None = None parent_id: str | None = None @dataclass class _SearchCall: # 1回の検索(query 単位)に対応する集約状態。 query: str = "" thinking: str = "" years: list[int] | None = None other_property: bool = False titles: list[str] = field(default_factory=list) class ChatStreamItem(NamedTuple): answer: str references: list[list[str]] queries: list[str] # オーケストレータが消費した累積 neoAI トークン。最終チャンクで確定する。 neoai_tokens: float # 最終回答に至るまでの途中経過パネル。 progress: list[ProgressPanel] def _build_panels(search_calls: dict[int, _SearchCall]) -> list[ProgressPanel]: panels: list[ProgressPanel] = [] if not search_calls: return panels panels.append( ProgressPanel(title="🔍 検索", content=f"{len(search_calls)} 件のクエリで検索しました", panel_id="search") ) for idx, call in sorted(search_calls.items()): lines: list[str] = [] if call.thinking.strip(): lines.append(call.thinking.strip()) scope = "他物件を含む" if call.other_property else "担当物件のみ" years = "、".join(f"{y}年度" for y in call.years) if call.years else "全年度" lines.append(f"**絞り込み**: {scope} / {years}") if call.titles: lines.append("**取得資料**: " + "、".join(call.titles)) else: lines.append("(該当資料なし)") panels.append( ProgressPanel( title=f"検索クエリ: {call.query}", content="\n".join(lines), panel_id=f"search-{idx}", parent_id="search", ) ) return panels def _reference_title(dp: DataPoint) -> str: if dp.meeting_no == 0: header = f"{dp.year}年度{dp.meeting_type}{dp.doc_type}" else: header = f"{dp.year}年度第{dp.meeting_no}回{dp.meeting_type}{dp.doc_type}" parts = [header] if dp.agenda_title: parts.append(dp.agenda_title) if dp.attachment_title: parts.append(f"({dp.attachment_title})") return " ".join(parts) def _build_messages(history: list[tuple[str, str]], question: str, input_files: list[InputFile]) -> Messages: messages: list[Message] = [] for user, assistant in history: if user: messages.append(Message(role=Role.USER, content=Content(root=user))) if assistant: messages.append(Message(role=Role.ASSISTANT, content=Content(root=assistant))) messages.append(Message(role=Role.USER, content=Content(root=question), input_files=input_files)) return Messages(root=messages) class ChatUseCase: @inject def __init__(self, completion: ChatCompletionUseCase, file_processor: IFileProcessorService) -> None: self._completion = completion self._file_processor = file_processor def _build_input_files(self, paths: list[str]) -> list[InputFile]: files: list[InputFile] = [] for p in paths: path = Path(p) files.append( self._file_processor.process_input_file(path.read_bytes(), path.name, PdfParser.DOCUMENT_INTELLIGENCE) ) return files def stream( self, question: str, property: Property, history: list[tuple[str, str]], input_files: list[str], tools: list[Tool] | None = None, search_other_property: bool | None = None, ) -> Iterator[ChatStreamItem]: """エージェント(検索ツールループ)で回答を生成し、途中経過(検索)と回答をストリームする。""" domain_files = self._build_input_files(input_files or []) chat_input = ChatCompletionInput( messages=_build_messages(history, question, domain_files), index_name=CARINA_INDEX_NAME, my_property=property, search_other_property=search_other_property, input_files=domain_files or None, tools=tools, ) answer = "" references: list[list[str]] = [] queries: list[str] = [] neoai_tokens = 0.0 search_calls: dict[int, _SearchCall] = {} ref_labels: dict[str, str] = {} seen_queries: set[str] = set() search_idx = 0 for event in self._completion.create_agentic_chat_stream(chat_input): if isinstance(event, QueryGeneratedEvent): search_idx += 1 query = event.query or "" search_calls[search_idx] = _SearchCall( query=query, thinking=event.thinking, years=event.years, other_property=event.other_property, ) if query and query not in seen_queries: seen_queries.add(query) queries.append(query) elif isinstance(event, ChatCompletionDataPointsEvent): call = search_calls.get(search_idx) for dp in event.data_points: title = _reference_title(dp) if title not in ref_labels: ref_labels[title] = f"[{len(references) + 1}]" references.append([ref_labels[title], title]) if call is not None and title not in call.titles: call.titles.append(title) elif isinstance(event, ChatCompletionAnswerEvent): answer = event.answer elif isinstance(event, TokenUsageEvent): neoai_tokens = event.neoai_token else: continue yield ChatStreamItem(answer, references, queries, neoai_tokens, _build_panels(search_calls)) class CarinaDemoApp: # 担当物件は経堂・祖師谷のみ表示。表示ラベルは短縮するが、 # 内部値(検索インデックス/メタデータ照合に使う enum 値)は変更しない。 _PROPERTY_CHOICES: ClassVar[list[tuple[str, str]]] = [ ("経堂", Property.KYODO_PM.value), ("祖師谷", Property.SOSHIGAYA_OKURA_PHO.value), ] _FEEDBACK_CATEGORY_CHOICES: ClassVar[list[str]] = [c.value for c in FeedbackCategory] def __init__( self, user_repository: IUserRepository, log_repository: IChatLogRepository, ) -> None: self._user_repository = user_repository self._log_repository = log_repository @staticmethod def _utc_now_iso() -> str: return datetime.now(timezone.utc).isoformat(timespec="milliseconds").replace("+00:00", "Z") # ── assets ────────────────────────────────────────────────────────── @staticmethod def _custom_css() -> str: return (ASSETS_DIR / "style.css").read_text(encoding="utf-8") @staticmethod def _placeholder() -> str: return (ASSETS_DIR / "placeholder.html").read_text(encoding="utf-8") # ── helpers ───────────────────────────────────────────────────────── @staticmethod def _to_messages(conversation: list[tuple[str, str]]) -> list[dict[str, str]]: messages: list[dict[str, str]] = [] for question, answer in conversation: messages.append({"role": "user", "content": question}) messages.append({"role": "assistant", "content": answer}) return messages @staticmethod def _on_files_uploaded(new_files: list[str] | None, existing_files: list[str] | None) -> list[str]: return list(dict.fromkeys([*(existing_files or []), *(new_files or [])])) @staticmethod def _progress_messages(progress: list[ProgressPanel], status: str) -> list[dict[str, Any]]: messages: list[dict[str, Any]] = [] for panel in progress: metadata: dict[str, Any] = {"title": panel.title, "status": status} if panel.panel_id is not None: metadata["id"] = panel.panel_id if panel.parent_id is not None: metadata["parent_id"] = panel.parent_id messages.append({"role": "assistant", "content": panel.content, "metadata": metadata}) return messages @staticmethod def _render_attached_files(files: list[str], state: gr.State) -> None: if not files: return with gr.Column(elem_classes=["attached-files"]): for idx, file_path in enumerate(files): with gr.Row(elem_classes=["attached-file-row"]): gr.Markdown(f"{Path(file_path).name}") remove_btn = gr.Button( "×", size="sm", scale=0, min_width=40, elem_classes=["remove-file-btn"], ) remove_btn.click( fn=lambda current, i=idx: current[:i] + current[i + 1 :], inputs=[state], outputs=[state], ) # ── handlers ──────────────────────────────────────────────────────── def _check_ip_access(self, request: gr.Request) -> tuple[dict[str, Any], dict[str, Any]]: """ページ読込時に呼ばれ、許可リスト外のIPなら門前払いする。 戻り値は (login_page, denied_page) の表示更新。 HF Spaces 経由では実IPは x-forwarded-for の最左値で取得する。 XFF が無い場合(プロキシ非経由=ローカル開発など)は直結IPで判定する。 """ forwarded_for = request.headers.get("x-forwarded-for") if request else None client_ip = extract_client_ip(forwarded_for) if client_ip is None and request is not None and request.client is not None: client_ip = request.client.host self._record_access_log(client_ip) if is_ip_allowed(client_ip, load_allowed_ips()): return gr.update(visible=True), gr.update(visible=False) return gr.update(visible=False), gr.update(visible=True) def _record_access_log(self, client_ip: str | None) -> None: """アクセス元IPと日時を chat ログとは別フォルダに永続化する。失敗してもページ表示は止めない。""" try: self._log_repository.append_access( AccessLogEntry( timestamp=self._utc_now_iso(), client_ip=client_ip, ) ) except Exception as e: # logging must never break the UI print(f"[carina] failed to record access log: {e}") def _handle_login( self, username: str, password: str ) -> tuple[dict[str, Any], dict[str, Any], dict[str, Any], str, str, str]: if self._user_repository.verify(username, password): return ( gr.update(visible=False), # login_page gr.update(visible=True), # main_content gr.update(visible=False), # login_message username, uuid.uuid4().hex, # session_id uuid.uuid4().hex, # conversation_id ) return ( gr.update(visible=True), gr.update(visible=False), gr.update( value="
", visible=True, ), "", "", "", ) def _answer( self, question: str, conversation: list[tuple[str, str]], property_label: str, search_other_property: bool, input_files: list[str], username: str, session_id: str, conversation_id: str, ) -> Iterator[tuple[str, list[tuple[str, str]], Any, str]]: if not question.strip(): yield "", conversation, self._to_messages(conversation), "" return property = Property(property_label) # モデルは1種類のため既定モデルに固定(UIからの選択は廃止)。 model = MODEL_CHOICES[DEFAULT_MODEL_LABEL] # チェックON→他物件を強制。OFF→None でAIの自動判定に委ねる。 other_property_override = True if search_other_property else None turn_id = uuid.uuid4().hex base_messages = self._to_messages(conversation) answer = "" references: list[list[str]] = [] neoai_tokens = 0.0 progress: list[ProgressPanel] = [] injector = create_injector(CARINA_INDEX_NAME, model_name=model) chat = injector.get(ChatUseCase) for chunk in chat.stream( question, property, conversation, input_files, search_other_property=other_property_override, ): answer, references, _, neoai_tokens, progress = chunk streaming = [ *base_messages, {"role": "user", "content": question}, *self._progress_messages(progress, status="pending"), ] if answer: streaming.append({"role": "assistant", "content": answer}) yield ( "", conversation, streaming, turn_id, ) new_conversation = [*conversation, (question, answer)] self._record_chat_log( username, session_id, conversation_id, turn_id, question, answer, references, neoai_tokens ) yield ( "", new_conversation, [ *base_messages, {"role": "user", "content": question}, *self._progress_messages(progress, status="done"), {"role": "assistant", "content": answer}, ], turn_id, ) def _record_chat_log( self, username: str, session_id: str, conversation_id: str, turn_id: str, question: str, answer: str, references: list[list[str]], neoai_tokens: float, ) -> None: if not username or not session_id or not conversation_id: return try: self._log_repository.append_log( ChatLogEntry( timestamp=self._utc_now_iso(), username=username, session_id=session_id, conversation_id=conversation_id, conversation_turn_id=turn_id, question=question, answer=answer, references=[r[1] for r in references], neoai_token=neoai_tokens, ) ) except Exception as e: # logging must never break the UI print(f"[carina] failed to record chat log: {e}") @staticmethod def _on_clear() -> tuple[list[tuple[str, str]], list[dict[str, str]], list[str], str, str, None]: return ( [], # conversation_state [], # chatbot [], # input_files_state uuid.uuid4().hex, # conversation_id_state "", # feedback_textbox None, # feedback_category ) def _submit_feedback( self, comment: str, category: str | None, username: str, session_id: str, conversation_id: str, turn_id: str | None, ) -> tuple[str, str | None]: if not category: gr.Warning("フィードバック項目を選択してください") return comment, category if not comment.strip(): gr.Warning("コメントを入力してください") return comment, category if not username or not session_id or not conversation_id: gr.Warning("ログインしてください") return comment, category try: self._log_repository.append_feedback( FeedbackLog( feedback_id=uuid.uuid4().hex, timestamp=self._utc_now_iso(), username=username, session_id=session_id, conversation_id=conversation_id, conversation_turn_id=turn_id or None, category=FeedbackCategory(category), comment=comment.strip(), ) ) except Exception: gr.Warning("フィードバックの保存に失敗しました") return comment, category gr.Info("フィードバックを送信しました") return "", None # ── UI ────────────────────────────────────────────────────────────── def build(self) -> gr.Blocks: with gr.Blocks(title="Carina Demo") as demo: gr.HTML(f"") current_user = gr.State("") session_id_state = gr.State("") conversation_id_state = gr.State("") # 初期は非表示。demo.load の IP チェック後に login / denied のどちらかを出す # (初期 visible=True だとログイン画面が一瞬見えてから 403 に差し替わるため)。 with gr.Column(visible=False, elem_classes=["login-container"]) as login_page: gr.HTML("