> For the complete documentation index, see [llms.txt](https://genos-docs.gitbook.io/default/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://genos-docs.gitbook.io/default/v1.9.2/admin-management/settings/extension-service/ai-meeting-guide.md).

# AI 회의록 코드서빙 및 워크플로우 구현 가이드

AI 회의록은 회의 오디오를 업로드하면 화자별로 분리·전사·요약하는 기능입니다. 활성화를 위해 다음 3가지 구성요소를 GenOS에 배포해야 합니다.

<table><thead><tr><th width="187.56640625">구성요소</th><th>형태</th><th>책임</th></tr></thead><tbody><tr><td><strong>화자분리 코드서빙</strong></td><td>Python 코드서빙 (<code>service.py</code>)</td><td>오디오 파일을 받아 화자별 발화 구간(<code>start</code>, <code>end</code>, <code>speaker</code>) 추출</td></tr><tr><td><strong>STT 워크플로우</strong></td><td>Python 워크플로우 (<code>llmops-workflow-api-audio:latest</code>)</td><td>세그먼트별로 오디오를 잘라 STT 모델 호출, 텍스트 스트리밍 반환</td></tr><tr><td><strong>요약 워크플로우</strong></td><td>에이전트 플로우 (LLM 노드)</td><td>화자·시각이 포함된 전사 텍스트를 받아 회의록 요약 생성</td></tr></tbody></table>

각 구성요소는 **AI 회의록 Chat API와의 계약(스키마)** 만 지키면 됩니다. 내부 구현(모델 선택, 최적화, 환경변수 등)은 자유입니다.

***

## 1. 화자분리 코드서빙

### 1-1. 입출력 스키마

#### 호출 방식

Chat API는 코드서빙 엔드포인트에 단발성 HTTP POST를 보냅니다.

```
POST {code_serving_url}/json
Content-Type: application/json
```

#### 입력 스키마

`async def service(config, data)`의 `data` 객체로 다음 필드가 주입됩니다.

<table><thead><tr><th width="112.171875">필드</th><th width="73.171875">타입</th><th width="79.53125">필수</th><th>설명</th></tr></thead><tbody><tr><td><code>audio_url</code></td><td>string</td><td>✅</td><td>오디오 파일 URL. Chat API는 MinIO presigned URL(클러스터 내부 ClusterIP 기반)을 전달하며, 저장 파일은 항상 .m4a 포맷. 코드서빙 내부에서 디코딩 책임 (참조 구현은 torchaudio + ffmpeg backend 사용)</td></tr></tbody></table>

`config` 객체는 코드서빙 플랫폼이 주입하는 런타임 설정이며, 본 호출 계약에는 포함되지 않습니다.

#### 출력 스키마

HTTP 200 응답으로 다음 JSON을 반환합니다.

<table><thead><tr><th width="115.28515625">필드</th><th width="75.59765625">타입</th><th width="75.2265625">필수</th><th>설명</th></tr></thead><tbody><tr><td><code>segments</code></td><td>array</td><td>✅</td><td>화자 발화 구간 목록. 누락 시 빈 배열로 간주되어 화자분리 실패 처리</td></tr></tbody></table>

**`segments[]` 원소**

<table><thead><tr><th width="117.34375">필드</th><th width="166.5703125">타입</th><th>단위 / 규칙</th><th>설명</th></tr></thead><tbody><tr><td><code>start</code></td><td>number(float)</td><td>초(seconds)</td><td>발화 구간 시작 시각</td></tr><tr><td><code>end</code></td><td>number(float)</td><td>초(seconds), <code>end > start</code></td><td>발화 구간 종료 시각</td></tr><tr><td><code>speaker</code></td><td>integer</td><td>1부터 시작하는 양의 정수</td><td>화자 식별자. 동일 인물은 동일 번호</td></tr></tbody></table>

#### 제약 사항

* `segments`는 `start` 오름차순으로 정렬되어야 합니다 (Chat API가 정렬을 가정하고 `sequence` 부여).
* `speaker`는 **반드시 정수**여야 합니다 (문자열·실수 불가).
* 발화가 없으면 `{"segments": []}` 반환. `null`·필드 누락 금지.

#### 응답 예시

```json
{
  "segments": [
    { "start": 0.512,  "end": 3.847,  "speaker": 1 },
    { "start": 3.912,  "end": 7.231,  "speaker": 2 },
    { "start": 7.305,  "end": 12.480, "speaker": 1 }
  ]
}
```

#### Chat API의 후속 처리 (참고)

1. `unique_speaker_indices = sorted({d["speaker"] for d in segments})` — 화자 마스터(`Speaker 1`, `Speaker 2`, …) 생성
2. 등장 순서대로 `sequence` 부여 → STT 워크플로우의 `segments[]` 입력으로 변환

### 1-2. 최소 구현 예시

계약만 만족하는 가장 단순한 형태입니다. 실제 동작 가능한 코드는 §1-4 참조 구현을 보세요.

```python
async def service(config, data):
    audio_url = data["audio_url"]

    # 1) audio_url 다운로드
    # 2) 임의의 화자분리 모델 추론
    # 3) 결과를 (start, end, speaker) 형식으로 변환
    raw_segments = your_diarization_model(audio_url)

    segments = [
        {
            "start": round(s.start, 3),         # 필수: float, 초
            "end": round(s.end, 3),             # 필수: float, 초
            "speaker": int(s.speaker_label),    # 필수: 1부터 시작 정수
        }
        for s in raw_segments
    ]
    return {"segments": segments}
```

### 1-3. 구현 절차

1. 코드서빙 생성
2. 해당 코드서빙을 코드스페이스에서 git clone
3. 프로젝트 루트에 `models` 폴더를 만들어 화자분리 모델을 업로드 또는 다운로드
4. 패키지 요구사항을 프로젝트 루트 `requirements.txt`에 작성

   작성 예시:

   ```
   pyannote.audio==3.4.0
   torchaudio==2.6.0
   soundfile
   onnxruntime-gpu
   numpy
   ```

   * 해당 패키지들과 의존성 패키지들은 PyPi 패키지 메뉴에 업로드되어 있어야 합니다.
   * 혹은 프로젝트 루트에 `packages` 폴더를 만들어 설치할 패키지 whl 파일을 포함시키면 PyPi 서버 대신 해당 경로에서 패키지를 설치할 수 있습니다.
5. `service.py` 작성 (§1-2 최소 구현 또는 §1-4 참조 구현)

### 1-4. 참조 구현 (pyannote.audio 기반)

> 아래는 운영 환경에서 검증된 참조 구현입니다. 다음 항목은 **구현 자유 영역**이며 다른 모델·라이브러리 사용 시 변경 가능합니다:
>
> * 화자분리 모델 (pyannote 외 NeMo, SpeechBrain 등)
> * GPU 최적화 (FP16 autocast, warmup)
> * 동시 호출 제어 (Semaphore, 단일 GPU 기준 N=1)
> * 무음 패딩 (`pad_seconds=11.0`은 pyannote 임베딩 윈도우 회피용)
> * ffmpeg backend 패치 (pyannote 한정 안전망)
> * 로깅 포맷, 임시 파일 처리, GPU 메모리 정리
>
> **§1-1 계약(`audio_url` 입력, `segments[{start, end, speaker}]` 출력)만 어기지 않으면 됩니다.**

* [**service.py**](http://service.py) **전체 코드 펼치기**

  ```python
  import os
  import sys
  import time
  import asyncio
  import logging
  import tempfile
  import warnings
  import torch
  import torchaudio
  from typing import Any, Dict

  logger = logging.getLogger(__name__)

  # pyannote의 짧은 발화 윈도우에서 std(correction=1) 호출 시 발생하는 cosmetic 경고 차단.
  # 결과에 영향 없는 알려진 노이즈라 로그에서 제거.
  warnings.filterwarnings("ignore", category=UserWarning, module="pyannote")

  # PyTorch 2.6+ weights_only=True 기본값 비활성화 (pyannote 모델 호환)
  _original_load = torch.load
  torch.load = lambda *args, **kwargs: _original_load(*args, **{**kwargs, 'weights_only': False})

  from pyannote.audio import Pipeline

  # ─────────────── M4A 등 비-WAV 포맷 지원 ───────────────
  # libsndfile(soundfile) backend는 WAV/FLAC/OGG만 지원하므로,
  # pyannote 내부에서 파일로부터 로드할 경우를 대비해 ffmpeg backend를 기본값으로 override.
  # (본 서비스는 in-memory waveform을 직접 넘겨서 이 경로를 거의 타지 않지만 안전망으로 유지.)
  from pyannote.audio.core.io import Audio as _PyannoteAudio

  _original_audio_init = _PyannoteAudio.__init__

  def _patched_audio_init(self, sample_rate=None, mono=None, backend=None):
      _original_audio_init(
          self,
          sample_rate=sample_rate,
          mono=mono,
          backend=backend if backend is not None else "ffmpeg",
      )

  _PyannoteAudio.__init__ = _patched_audio_init
  # ────────────────────────────────────────────────────────

  _BASE_DIR = os.path.dirname(os.path.abspath(__file__))
  _MODEL_DIR = os.path.join(_BASE_DIR, "models")

  _CONFIG = f"""version: 3.1.0

  pipeline:
    name: pyannote.audio.pipelines.SpeakerDiarization
    params:
      clustering: AgglomerativeClustering
      embedding: {_MODEL_DIR}/pyannote_embedding_wespeaker.bin
      embedding_batch_size: 32
      embedding_exclude_overlap: true
      segmentation: {_MODEL_DIR}/segmentation-3.0/pytorch_model.bin
      segmentation_batch_size: 32

  params:
    clustering:
      method: centroid
      min_cluster_size: 12
      threshold: 0.7045654963945799
    segmentation:
      min_duration_off: 0.0
  """

  def _load_pipeline() -> Pipeline:
      tmp = tempfile.NamedTemporaryFile(mode="w", suffix=".yaml", delete=False)
      tmp.write(_CONFIG)
      tmp.close()
      try:
          pipeline = Pipeline.from_pretrained(tmp.name)
      finally:
          os.unlink(tmp.name)

      device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
      pipeline.to(device)
      # NOTE: pyannote의 .to()는 dtype을 안 받음. FP16은 inference 시점에
      # torch.autocast로 mixed precision 적용 (_run_diarization 참조).
      return pipeline

  _pipeline = _load_pipeline()

  # 동시 호출 제한 — pyannote는 thread-safe 하지 않고, 단일 GPU에서 동시 호출은
  # activation 메모리 누적으로 OOM·race condition을 유발한다. 단일 GPU 환경 기준 N=1.
  _concurrency = asyncio.Semaphore(1)

  def _run_diarization(audio_input: dict):
      """파이프라인 추론 실행. CUDA 환경에선 autocast로 FP16 mixed precision 적용.

      pyannote.Pipeline.to()는 device만 받아 dtype 캐스팅 불가. autocast를 쓰면
      가중치는 FP32 유지하면서 forward 활성화만 FP16으로 처리 (PyTorch 표준 mixed precision).
      """
      if torch.cuda.is_available():
          with torch.autocast(device_type="cuda", dtype=torch.float16):
              return _pipeline(audio_input)
      return _pipeline(audio_input)

  def _warmup() -> None:
      """Startup에서 dummy 입력으로 파이프라인 1회 실행 (autocast 포함).

      첫 실 요청에서 발생하는 CUDA 커널 JIT 컴파일·cuDNN 알고리즘 선택·가중치 페이지 로드를
      미리 수행. 실패 시 모델·GPU 환경이 망가진 상태이므로 즉시 종료해 오케스트레이터에
      재시작을 위임한다.
      """
      sample_rate = 16000
      dummy = torch.zeros(1, sample_rate)  # 1초 무음
      _run_diarization({"waveform": dummy, "sample_rate": sample_rate})

  try:
      _warmup()
  except Exception as e:
      logger.error(f"Pipeline warmup failed: {e}", exc_info=True)
      sys.exit(1)

  def _load_and_pad(path: str, pad_seconds: float = 11.0) -> tuple[torch.Tensor, int]:
      """파일을 메모리에 로드하고 끝에 무음 패딩을 추가해 (waveform, sample_rate) 반환.

      pyannote는 임베딩 추출 시 10초 고정 윈도우를 사용하는데, 마지막 발화가 파일 끝 근처에
      있으면 윈도우가 짧아져 torch.vstack 크기 불일치가 발생한다. 무음 패딩으로 방지.

      디스크 재쓰기 없이 메모리에서만 패딩해 pyannote에 dict로 직접 전달한다
      (파일 재저장 시 lossy 코덱 재인코딩으로 샘플 정확성이 깨지는 문제 회피).
      """
      waveform, sr = torchaudio.load(path, backend="ffmpeg")
      pad_samples = int(pad_seconds * sr)
      silence = torch.zeros(waveform.shape[0], pad_samples)
      padded = torch.cat([waveform, silence], dim=1)
      return padded, sr

  def _semaphore_state() -> tuple[int, int]:
      """현재 _concurrency Semaphore의 (running, waiting) 반환."""
      capacity = 1  # _concurrency = Semaphore(1)
      free = _concurrency._value
      running = capacity - free
      waiting = len(_concurrency._waiters) if _concurrency._waiters else 0
      return running, waiting

  async def service(config: Dict[str, Any], data: Dict[str, Any]):
      from urllib.parse import urlparse

      audio_url = data.get("audio_url")
      if not audio_url:
          return {"error": "audio_url is required"}

      # 식별자 — audio_url의 파일명에서 meeting_id 추출 (예: 1234.m4a → 1234)
      parsed = urlparse(audio_url)
      basename = os.path.basename(parsed.path) if parsed.path else "unknown"
      request_id = basename.rsplit(".", 1)[0] if "." in basename else basename

      ext = os.path.splitext(parsed.path)[1] or ".wav"

      with tempfile.NamedTemporaryFile(suffix=ext, delete=False) as tmp:
          tmp_path = tmp.name

      # 진입 로그
      running, waiting = _semaphore_state()
      logger.info(
          f"[Diarization][{request_id}] === 요청 진입 === "
          f"running={running}/1 waiting={waiting}"
      )
      request_start = time.monotonic()

      try:
          # ─── 다운로드 (Semaphore 밖, 동시 다운로드 허용) ─────────────────
          if parsed.scheme == "file":
              import shutil
              await asyncio.to_thread(shutil.copy, parsed.path, tmp_path)
          else:
              import urllib.request
              await asyncio.to_thread(urllib.request.urlretrieve, audio_url, tmp_path)

          file_size_mb = os.path.getsize(tmp_path) / (1024 * 1024)

          # ─── Semaphore 진입 ──────────────────────────────────────────────
          wait_start = time.monotonic()
          async with _concurrency:
              queue_wait = time.monotonic() - wait_start
              running, waiting = _semaphore_state()
              logger.info(
                  f"[Diarization][{request_id}] 슬롯 획득 "
                  f"running={running}/1 waiting={waiting} queue_wait={queue_wait:.1f}s"
              )

              # 오디오 로드 + 길이 측정
              waveform, sr = await asyncio.to_thread(_load_and_pad, tmp_path)
              # waveform shape: (channels, samples) — pad_seconds(11s)가 끝에 추가됐으므로 빼줌
              audio_duration_sec = (waveform.shape[-1] / sr) - 11.0

              logger.info(
                  f"[Diarization][{request_id}] 추론 시작 "
                  f"audio={audio_duration_sec:.1f}s size={file_size_mb:.1f}MB"
              )

              # 추론 (FP16 autocast)
              inference_start = time.monotonic()
              try:
                  diarization = await asyncio.to_thread(
                      _run_diarization,
                      {"waveform": waveform, "sample_rate": sr},
                  )
              except torch.cuda.OutOfMemoryError as e:
                  elapsed = time.monotonic() - inference_start
                  logger.error(
                      f"[Diarization][{request_id}] CUDA OOM after {elapsed:.1f}s: {e}"
                  )
                  return {"error": "out_of_memory"}
              finally:
                  del waveform

              inference_elapsed = time.monotonic() - inference_start

          # ─── 결과 가공 (Semaphore 밖) ────────────────────────────────────
          speaker_map = {}
          speaker_counter = 1
          segments = []

          for turn, _, speaker in diarization.itertracks(yield_label=True):
              if speaker not in speaker_map:
                  speaker_map[speaker] = speaker_counter
                  speaker_counter += 1
              segments.append({
                  "start": round(turn.start, 3),
                  "end": round(turn.end, 3),
                  "speaker": speaker_map[speaker],
              })

          # 종합 요약
          rtf = inference_elapsed / audio_duration_sec if audio_duration_sec > 0 else 0
          total_elapsed = time.monotonic() - request_start
          logger.info(
              f"[Diarization][{request_id}] === 완료 === "
              f"audio={audio_duration_sec:.1f}s size={file_size_mb:.1f}MB "
              f"inference={inference_elapsed:.1f}s rtf={rtf:.3f} "
              f"total={total_elapsed:.1f}s "
              f"speakers={len(speaker_map)} segments={len(segments)}"
          )

          return {"segments": segments}
      finally:
          # 호출 간 GPU 메모리 누수 방지: PyTorch caching allocator가 잡고 있는 미사용 블록 회수.
          if torch.cuda.is_available():
              torch.cuda.empty_cache()
          os.unlink(tmp_path)
  ```

***

## 2. STT 워크플로우

### 2-1. 입출력 스키마

#### 호출 방식

Chat API는 워크플로우를 **스트리밍 모드**로 호출합니다.

```
POST {gateway_url}/llmops/workflow/{workflow_id}/run/v2
Content-Type: application/json
```

#### 입력 스키마

`async def run(data)`의 `data` 객체로 다음 필드가 주입됩니다.

<table><thead><tr><th width="150.06640625">필드</th><th width="102.62890625">타입</th><th width="76.078125">필수</th><th>설명</th></tr></thead><tbody><tr><td><code>bucket</code></td><td>string</td><td>✅</td><td>오디오 파일이 저장된 MinIO 버킷명</td></tr><tr><td><code>object_key</code></td><td>string</td><td>✅</td><td>MinIO 오브젝트 키. Chat API는 업로드된 모든 오디오를 .m4a 포맷으로 변환해 MinIO에 저장하므로 항상 .m4a 확장자로 전달됨. 워크플로우 내부에서는 ffmpeg 등으로 디코딩 후 STT 모델 입력 포맷(보통 mono 16kHz WAV)으로 변환</td></tr><tr><td><code>segments</code></td><td>array</td><td>✅</td><td>STT 대상 세그먼트 목록 (화자분리 결과를 Chat API가 변환해 주입)</td></tr><tr><td><code>stream</code></td><td>boolean</td><td>✅</td><td>항상 <code>true</code></td></tr></tbody></table>

**`segments[]` 원소 — 필수 필드**

| 필드         | 타입            | 필수 | 설명                                |
| ---------- | ------------- | -- | --------------------------------- |
| `sequence` | integer       | ✅  | 세그먼트 식별자(1부터). 응답 매칭 키            |
| `start`    | number(float) | ✅  | 원본 오디오 기준 시작 시각(초)                |
| `end`      | number(float) | ✅  | 원본 오디오 기준 종료 시각(초), `end > start` |

#### 출력 스키마 — yield 이벤트 규약

워크플로우 Python step은 **async generator**로 다음 형식의 dict를 yield 해야 합니다. 모든 yield는 GenOS 워크플로우 런타임 공통 규약 `{"event": <string>, "data": <any>}` envelope을 따라야 하며, 런타임이 자동으로 `data: <json>\\n\\n` SSE 라인으로 변환합니다.

**Chat API가 인식하는 이벤트**

| 이벤트         | 발행 시점              | `data` 필수 필드                    | 처리                      |
| ----------- | ------------------ | ------------------------------- | ----------------------- |
| `utterance` | 세그먼트 1개 STT 완료 시마다 | `sequence`(int), `text`(string) | sequence로 매칭해 발화 텍스트 저장 |
| `result`    | 모든 세그먼트 성공 완료 후    | (없음, `{}` 허용)                   | 워크플로우 종료 신호             |

Chat API는 `utterance`·`result` 외 이벤트는 무시합니다. 로깅·디버깅용 token 등을 자유롭게 yield해도 무방하나 회의록에 반영되지 않습니다.

#### 제약 사항

* utterance 이벤트는 입력 segments의 모든 원소에 대해 정확히 1회씩 발행되어야 합니다. 누락 시 Chat API가 "분석 중 오류" 메시지로 전체 실패 처리합니다.
* utterance는 완료 순서대로 yield 해도 됩니다 (Chat API가 sequence 기준으로 재정렬).
* 세그먼트 단위 실패 처리는 워크플로우 작성자 재량입니다:
  * 옵션 A: text="" 으로 utterance를 발행 → Chat API는 빈 텍스트도 정상 발화로 카운트하여 회의 단위 성공 처리. 빈 값은 FE에 그대로 전달됨.
  * 옵션 B: 예외 raise → 워크플로우 HTTP 에러 → Chat API가 회의 단위 전체 실패 처리.
* result 이벤트가 한 번도 발행되지 않으면 Chat API가 STT 실패로 간주합니다.
* 워크플로우 자체적으로 "세그먼트 실패" 이벤트를 정의해도 Chat API가 인식하지 못합니다 (계약상 utterance·result만 처리).

#### 이벤트 예시 (워크플로우 → Chat API SSE 와이어)

```
data: {"event": "utterance", "data": {"sequence": 2, "text": "안녕하세요"}}

data: {"event": "utterance", "data": {"sequence": 1, "text": "네 시작하시죠"}}

data: {"event": "result", "data": {}}
```

### 2-2. 최소 구현 예시

```python
async def run(data):
    bucket = data["bucket"]
    object_key = data["object_key"]
    segments = data["segments"]

    # 1) MinIO에서 오디오 다운로드
    audio_bytes = your_minio_download(bucket, object_key)

    # 2) 각 세그먼트별로 STT 호출
    for seg in segments:
        text = await your_stt_call(audio_bytes, seg["start"], seg["end"])
        yield {"event": "utterance", "data": {"sequence": seg["sequence"], "text": text}}

    # 3) 종료 신호
    yield {"event": "result", "data": {}}
```

### 2-3. 구현 절차

1. 워크플로우 생성
2. 도커 이미지: `llmops-workflow-api-audio:latest` (ffmpeg 내장)
3. Python step에 `run(data)` async generator 작성
4. STT 모델 서빙 ID를 환경변수로 주입 (예시 코드는 `STT_SERVING_ID`)
5. 필요 시 동시성·재시도·타임아웃 환경변수 설정

### 2-4. 참조 구현 (ffmpeg + httpx + 동시성 제어)

> 아래는 운영 환경에서 검증된 참조 구현입니다. 다음 항목은 **구현 자유 영역**입니다:
>
> * 디코딩 방식 (ffmpeg subprocess / 라이브러리)
> * 세그먼트 슬라이싱 방식 (메모리 / 디스크)
> * 동시 처리 한도·재시도·타임아웃 환경변수 (`STT_SEGMENT_CONCURRENCY`, `STT_MAX_RETRIES`, `STT_READ_TIMEOUT`, `STT_TASK_STAGGER_MS`)
> * STT 모델 서빙 선택 방식 (`STT_SERVING_ID` 환경변수 vs 하드코딩)
> * 로깅 포맷, 임시 파일 처리
>
> **§2-1 계약(입력 필드, `{event, data}` envelope, `utterance`·`result` 발행 규칙)만 어기지 않으면 됩니다.**

* **워크플로우 Python step 전체 코드 펼치기**

  ```python
  import asyncio
  import io
  import json
  import os
  import subprocess
  import sys
  import tempfile
  import wave
  from datetime import datetime

  def _log(msg: str) -> None:
      """워크플로우 stdout으로 직접 출력 (logger 핸들러 우회)."""
      ts = datetime.now().strftime("%H:%M:%S.%f")[:-3]
      print(f"[STT][{ts}] {msg}", flush=True, file=sys.stdout)

  # ─── 동시 세그먼트 처리 한도 (전역) ─────────────────────────────────────────
  # 환경변수 STT_SEGMENT_CONCURRENCY 로 무중단 튜닝.
  _SEGMENT_CONCURRENCY = int(os.environ.get("STT_SEGMENT_CONCURRENCY", "16"))
  _segment_sem = asyncio.Semaphore(_SEGMENT_CONCURRENCY)

  # ─── STT 호출 재시도·타임아웃 ──────────────────────────────────────────────
  _STT_MAX_RETRIES = int(os.environ.get("STT_MAX_RETRIES", "3"))
  _STT_READ_TIMEOUT = float(os.environ.get("STT_READ_TIMEOUT", "90"))

  # ─── 태스크 생성 시 미세 stagger ────────────────────────────────────────────
  # 초기 task 들이 동시에 STT 폭주하는 걸 막기 위한 짧은 간격 (ms 단위).
  _TASK_STAGGER_MS = float(os.environ.get("STT_TASK_STAGGER_MS", "2"))

  def _decode_to_wav_file(input_path: str) -> str:
      """원본 오디오 → mono 16kHz PCM WAV 로 1회 디코딩.

      ffmpeg subprocess 1회 호출. 결과 WAV 는 같은 디렉터리에 `.wav` suffix 로 생성.
      Whisper 등 STT 표준 입력 포맷(mono 16kHz)에 맞춰 다운샘플링 동시 수행.
      """
      output_path = input_path + ".wav"
      subprocess.run(
          [
              "ffmpeg", "-y", "-loglevel", "error",
              "-i", input_path,
              "-ar", "16000", "-ac", "1",
              "-c:a", "pcm_s16le",
              output_path,
          ],
          check=True,
      )
      return output_path

  def _slice_wav_to_bytes(wav_path: str, start_sec: float, duration_sec: float) -> io.BytesIO:
      """WAV 파일에서 구간만 seek+read → 새 WAV bytes 로 반환.

      stdlib `wave` 모듈만 사용. 외부 라이브러리·ffmpeg 미호출.
      WAV 가 uncompressed + frame-aligned 이라 lseek 한 번에 정확한 구간 추출 가능.
      """
      with wave.open(wav_path, "rb") as wf:
          sr = wf.getframerate()
          channels = wf.getnchannels()
          sample_width = wf.getsampwidth()
          start_frame = int(start_sec * sr)
          n_frames = int(duration_sec * sr)
          wf.setpos(start_frame)
          frames = wf.readframes(n_frames)

      buf = io.BytesIO()
      with wave.open(buf, "wb") as out:
          out.setnchannels(channels)
          out.setsampwidth(sample_width)
          out.setframerate(sr)
          out.writeframes(frames)
      buf.seek(0)
      return buf

  async def run(data):
      import httpx
      from minio import Minio

      bucket: str = data["bucket"]
      object_key: str = data["object_key"]
      segments: list = data["segments"]
      meeting_id = object_key.rsplit(".", 1)[0] if "." in object_key else object_key

      _log(
          f"[{meeting_id}] === 처리 시작 segments={len(segments)} "
          f"bucket={bucket} concurrency={_SEGMENT_CONCURRENCY} ==="
      )

      # ─── MinIO → 임시 파일 ────────────────────────────────────────────────
      minio_client = Minio(
          endpoint=os.environ["MINIO_ENDPOINT"],
          access_key=os.environ["MINIO_ACCESS_KEY"],
          secret_key=os.environ["MINIO_SECRET_KEY"],
          secure=os.environ.get("MINIO_SECURE", "false").lower() == "true",
      )

      suffix = os.path.splitext(object_key)[1] or ".m4a"
      with tempfile.NamedTemporaryFile(suffix=suffix, delete=False) as tmp:
          tmp_path = tmp.name
          resp = minio_client.get_object(bucket, object_key)
          try:
              for chunk in resp.stream(64 * 1024):
                  tmp.write(chunk)
          finally:
              resp.close()
              resp.release_conn()

      try:
          file_size = os.path.getsize(tmp_path)
          _log(f"[{meeting_id}] 다운로드 완료 tmp={tmp_path} size={file_size}B")
      except OSError as e:
          _log(f"[{meeting_id}] [FATAL] 다운로드 후 tmp 파일 확인 실패: {e}")
          raise

      # ─── 원본 오디오 1회 디코딩 → mono 16kHz WAV ──────────────────────────
      # ffmpeg subprocess 가 blocking 이므로 별도 스레드에서 실행.
      wav_path: str | None = None
      try:
          wav_path = await asyncio.to_thread(_decode_to_wav_file, tmp_path)
          wav_size = os.path.getsize(wav_path)
          _log(f"[{meeting_id}] WAV 디코딩 완료 {wav_path} size={wav_size/1024/1024:.1f}MB")
      except subprocess.CalledProcessError as e:
          _log(f"[{meeting_id}] [FATAL] WAV 디코딩 실패: {e}")
          # tmp_path 정리 후 raise (wav_path 는 아직 없음)
          try:
              os.unlink(tmp_path)
          except OSError:
              pass
          raise

      # ─── STT 서빙 설정 ────────────────────────────────────────────────────
      serving_id = os.environ.get("STT_SERVING_ID", "")
      if not serving_id:
          raise ValueError("STT_SERVING_ID 환경변수가 설정되지 않았습니다.")
      gateway_url = "<http://llmops-gateway-api-service:8080>"
      stt_url = f"{gateway_url}/rep/serving/{serving_id}/v1/audio/transcriptions"

      # ─── 세그먼트별 처리 (모든 예외 재시도, 최종 실패 시 raise) ───────────
      async def process_segment(seg: dict, client: httpx.AsyncClient) -> tuple[int, str]:
          seq = seg["sequence"]
          last_error: Exception | None = None

          # 슬롯 한 번 잡은 채로 모든 재시도 — 재시도가 큐 뒤로 밀리지 않음
          async with _segment_sem:
              if not os.path.exists(wav_path):
                  raise FileNotFoundError(f"seq={seq} WAV 파일 사라짐: {wav_path}")

              for attempt in range(_STT_MAX_RETRIES + 1):
                  try:
                      duration = seg["end"] - seg["start"]
                      # 디스크 seek + read 만 — ffmpeg 호출 없음, 매우 빠름
                      buf = await asyncio.to_thread(
                          _slice_wav_to_bytes, wav_path, seg["start"], duration
                      )

                      resp = await client.post(
                          stt_url,
                          data={
                              "model": serving_id,
                              "response_format": "json",
                              "language": "ko",
                              "temperature": "0",
                          },
                          files={"file": ("segment.wav", buf, "audio/wav")},
                      )
                      resp.raise_for_status()
                      text = resp.json().get("text", "")

                      if attempt > 0:
                          _log(f"[{meeting_id}] seq={seq} 재시도 성공 (attempt={attempt})")
                      return seq, text

                  except Exception as e:
                      last_error = e
                      if attempt < _STT_MAX_RETRIES:
                          _log(
                              f"[{meeting_id}] seq={seq} {type(e).__name__}: {e} "
                              f"재시도 예정 ({attempt + 1}/{_STT_MAX_RETRIES})"
                          )
                          # 슬롯 잡은 채로 backoff
                          await asyncio.sleep(3 * (2 ** attempt))   # 3s, 6s, 12s
                      else:
                          _log(
                              f"[{meeting_id}] seq={seq} 최종 실패 "
                              f"({_STT_MAX_RETRIES + 1}회 시도): {type(e).__name__}: {e}"
                          )

          raise RuntimeError(
              f"seq={seq} STT 처리 실패 ({_STT_MAX_RETRIES + 1}회 시도)"
          ) from last_error

      async def _spawn_tasks_with_stagger(sorted_segs, client):
          """초기 N개 task 가 동시에 STT 로 폭주하는 걸 막기 위해 미세 stagger 로 생성."""
          tasks_local: list[asyncio.Task] = []
          stagger_sec = _TASK_STAGGER_MS / 1000.0
          for i, seg in enumerate(sorted_segs):
              tasks_local.append(asyncio.create_task(process_segment(seg, client)))
              # 첫 _SEGMENT_CONCURRENCY 개만 stagger — 그 이후는 어차피 세마포어 큐
              if stagger_sec > 0 and i < _SEGMENT_CONCURRENCY:
                  await asyncio.sleep(stagger_sec)
          return tasks_local

      # ─── 모든 세그먼트 task 로 띄우고 완료 순서대로 yield ───────────────────
      tasks: list[asyncio.Task] = []
      try:
          async with httpx.AsyncClient(
              timeout=httpx.Timeout(
                  connect=10.0,   # 부하 상황 대비 여유
                  read=_STT_READ_TIMEOUT,
                  write=10.0,
                  pool=10.0,
              ),
              limits=httpx.Limits(
                  max_keepalive_connections=_SEGMENT_CONCURRENCY * 2,
                  max_connections=_SEGMENT_CONCURRENCY * 2,
              ),
          ) as client:
              sorted_segs = sorted(segments, key=lambda s: s["sequence"])
              total = len(sorted_segs)
              tasks = await _spawn_tasks_with_stagger(sorted_segs, client)
              _log(
                  f"[{meeting_id}] {total}개 task 생성 완료 "
                  f"(stagger={_TASK_STAGGER_MS}ms)"
              )

              completed = 0
              for finished in asyncio.as_completed(tasks):
                  try:
                      seq, text = await finished
                  except Exception as e:
                      # 한 세그먼트 실패 → 회의 전체 실패. 호출자에게 알리고 propagate.
                      _log(
                          f"[{meeting_id}] === 처리 중단 === "
                          f"세그먼트 실패로 전체 회의 실패: {type(e).__name__}: {e}"
                      )
                      yield json.dumps({
                          "event": "error",
                          "data": {
                              "message": f"STT 처리 실패: {e}",
                              "meeting_id": meeting_id,
                          },
                      }, ensure_ascii=False)
                      raise

                  completed += 1
                  free = _segment_sem._value
                  running = _SEGMENT_CONCURRENCY - free
                  waiting = len(_segment_sem._waiters) if _segment_sem._waiters else 0
                  _log(
                      f"[{meeting_id}] {completed}/{total} 완료 (seq={seq}) "
                      f"running={running}/{_SEGMENT_CONCURRENCY} waiting={waiting}"
                  )

                  yield json.dumps({
                      "event": "utterance",
                      "data": {"sequence": seq, "text": text},
                  }, ensure_ascii=False)

          _log(f"[{meeting_id}] === 처리 완료 total={total} ===")
          yield json.dumps({"event": "result", "data": {}}, ensure_ascii=False)
      finally:
          # 1) 남아있는 task 정리 — generator 가 일찍 종료되거나 예외 발생 시
          #    pending task 들이 좀비가 되는 걸 막는다. THIS 회의의 tasks 만 정리되므로
          #    다른 회의에서 돌고 있는 작업에 영향 없음.
          pending = [t for t in tasks if not t.done()]
          if pending:
              _log(f"[{meeting_id}] cleanup: cancelling {len(pending)} pending tasks")
              for t in pending:
                  t.cancel()
              await asyncio.gather(*tasks, return_exceptions=True)

          # 2) 임시 파일 정리 — 원본 m4a + 디코딩된 WAV 모두 삭제
          for p in filter(None, (tmp_path, wav_path)):
              try:
                  os.unlink(p)
              except OSError:
                  pass
  ```

***

## 3. 요약 워크플로우

### 3-1. 입출력 스키마

#### 호출 방식

Chat API는 STT 워크플로우와 동일하게 워크플로우를 **스트리밍 모드**로 호출합니다.

```
POST {gateway_url}/llmops/workflow/{workflow_id}/run/v2
Content-Type: application/json
```

#### 입력 스키마

| 필드         | 타입      | 필수 | 설명                                   |
| ---------- | ------- | -- | ------------------------------------ |
| `question` | string  | ✅  | 회의 전사 텍스트. 화자명·타임스탬프가 포함된 형식 (아래 참조) |
| `stream`   | boolean | ✅  | 항상 `true`                            |

**`question` 텍스트 형식**

Chat API가 utterance 테이블에서 다음 형식으로 조립해 전달합니다.

```
[0.5s] Speaker 1: 안녕하세요, 회의 시작하겠습니다.
[3.9s] Speaker 2: 네 시작하시죠.
[7.3s] Speaker 1: 오늘 안건은 ...
...
```

* `[<start_time>s]`: 발화 시작 시각 (start\_time이 없으면 생략)
* `<speaker.display_name>`: 화자 표시명 (기본 `Speaker 1`, `Speaker 2`, … / 사용자가 변경 가능)
* `<content>`: STT 텍스트

#### 출력 스키마 — yield 이벤트 규약

요약 워크플로우는 LLM 토큰 스트리밍과 최종 결과를 다음 2종 이벤트로 발행합니다.

| 이벤트      | 발행 시점            | `data` 필수 필드   | 처리                           |
| -------- | ---------------- | -------------- | ---------------------------- |
| `token`  | LLM 토큰 1개 생성 시마다 | string         | FE에 `summary_token` 이벤트로 재발행 |
| `result` | 생성 완료 후          | `text`(string) | 회의록 DB에 요약 본문 저장             |

#### 제약 사항

* `result` 이벤트의 `data.text` 필드는 **필수**입니다 ([chat-api/src/service/ai\_meeting/ai\_meeting\_processor.py:385](https://claude.ai/epitaxy/chat-api/src/service/ai_meeting/ai_meeting_processor.py:385) — `data["text"]`만 참조).
* `result` 이벤트가 한 번도 발행되지 않거나 `text`가 빈 문자열이면 Chat API가 "요약 생성 중 오류가 발생했습니다" 메시지로 사용자에게 알립니다.
* `token` 스트리밍 없이 `result`만 한 번에 발행해도 무방하나, FE에서 실시간 토큰 출력이 끊겨 UX가 저하됩니다.

#### 이벤트 예시 (워크플로우 → Chat API SSE 와이어)

```
data: {"event": "token", "data": "## 회의"}

data: {"event": "token", "data": " 핵심 요약\\n\\n"}

data: {"event": "token", "data": "프로젝트 일정 ..."}

data: {"event": "result", "data": {"text": "## 회의 핵심 요약\\n\\n프로젝트 일정 ..."}}
```

### 3-2. 구현 방법 — 에이전트 플로우

요약 워크플로우는 **GenOS Agentflow 3.0.0** 으로 구성합니다. Python step을 작성하지 않고 LLM 노드 1개로 충분하며, LLM 노드가 `token`·`result` 이벤트 발행을 자동 처리하므로 워크플로우 작성자는 **프롬프트와 모델 설정만** 정의하면 됩니다.

#### 노드 구성

```
[Start] ──> [LLM 0]
```

* **Start 노드**: 입력 진입점. 별도 설정 없음. Chat API가 보내는 `question` 필드가 그대로 다음 노드로 전달됨.
* **LLM 0 노드**: LLM 호출 및 토큰 스트리밍. 아래 설정 참조.

#### LLM 노드 설정

| 항목                     | 설정값                         | 비고                                                                               |
| ---------------------- | --------------------------- | -------------------------------------------------------------------------------- |
| **Model**              | `ChatMnc` 외 사내 LLM 또는 외부 모델 | 한국어 회의록 요약 품질이 검증된 모델 권장                                                         |
| **Max Tokens**         | `-1` (예시 기준)                | `-1`은 모델 최대 출력 한도 사용. 최대 90분의 장시간 회의록도 잘림 없이 요약. 짧은 회의만 처리한다면 `512~2048`로 제한해도 됨 |
| **Temperature**        | `0.9` (예시 기준)               | 요약의 정확성을 우선한다면 `0.1~0.3` 권장. 다양한 표현을 원하면 예시값 유지                                  |
| **Streaming**          | ✅ ON (필수)                   | **반드시 활성화**. 비활성화 시 `token` 이벤트가 발행되지 않아 FE에서 실시간 출력이 끊김                         |
| **Enable Memory**      | OFF                         | 회의별 1회성 호출이라 컨텍스트 메모리 불필요                                                        |
| **Return Response As** | `User Message`              | LLM 응답을 다음 노드의 user message로 반환 (단일 LLM 구성에선 그대로 `result.text`로 매핑됨)             |

#### Messages 설정

LLM 노드의 `Messages` 섹션에 **System Message 1개**를 추가합니다.

| 항목          | 값               |
| ----------- | --------------- |
| **Role**    | `System`        |
| **Content** | §3-3 프롬프트 예시 참조 |

> Content 필드에 변수 바인딩 표기(`{{question}}` 등)를 명시적으로 쓰지 않아도, Flowise 에이전트 플로우는 Start 노드로 들어온 입력을 자동으로 user message로 LLM에 전달합니다. System Message에는 **요약 지침만** 작성하면 됩니다.

### 3-3. 프롬프트 예시

```
당신은 전문적인 회의록 요약가입니다.

아래 회의 녹취록을 읽고 다음 형식으로 요약하세요.

[중요 원칙]
- 반드시 녹취록에 실제로 언급된 내용만 작성하세요.
- 녹취록에 없는 내용을 추측하거나 추가하지 마세요.
- 내용이 불분명하거나 언급되지 않은 항목은 "없음"으로 표시하세요.

## 회의 핵심 요약
(녹취록 전체 내용을 3~5줄로 요약. 언급된 내용만 포함)

## 주요 논의 사항
(녹취록에서 실제로 논의된 안건을 bullet point로 정리. 없으면 "없음")

## 결정 사항
(회의에서 명확히 결정된 내용만 작성. 없으면 "없음")

## 액션 아이템
(담당자와 할 일이 명확히 언급된 경우만 포함. 없으면 "없음")
```

> 프롬프트의 헤더 구조(`## 회의 핵심 요약` 등)는 FE 렌더링에 직접 영향을 주지 않습니다. 마크다운으로 출력되며, FE가 마크다운 렌더러로 그대로 표시합니다.

### 3-4. 구현 자유 영역

* LLM 모델 선택 (Claude / GPT / 사내 LLM 등)
* 프롬프트 문구·섹션 구조
* 온도(temperature)·max\_tokens 등 LLM 파라미터
* 전처리/후처리 노드 추가 (예: 발화자 매핑 보정, 마크다운 후처리)

단, **`question` 입력 필드명, `token`·`result` yield 이벤트 규약, `result.data.text` 필드**는 어겨선 안 됩니다.

***

## 4. 프론트엔드 수신 이벤트 (공통 참고)

Chat API는 워크플로우의 이벤트를 받아 **회의별 Redis 채널로 재발행**하며, FE는 SSE 구독으로 다음 이벤트를 수신합니다.

| 이벤트              | 발행 시점                           | `data` 필드                                              |
| ---------------- | ------------------------------- | ------------------------------------------------------ |
| `initial_data`   | 화자분리 완료 직후                      | `speakers[]`, `utterances[]` (skeleton, content 없음)    |
| `audio_ready`    | 오디오 업로드·화자분리 완료 후 STT 시작 직전     | 회의 메타데이터                                               |
| `utterance_text` | 워크플로우 `utterance` 이벤트 1건 수신 시마다 | `id`(utterance PK), `sequence`(int), `content`(string) |
| `summary_token`  | 요약 워크플로우 `token` 이벤트 수신 시마다     | `content`(string)                                      |
| `done`           | 단계 완료 시 (분석 완료 / 요약 완료 / 실패 종결) | `status`(AiMeetingStatus)                              |

> FE 개발자는 워크플로우의 `utterance`·`token` 이벤트가 아니라 Chat API가 가공한 `utterance_text`·`summary_token` 이벤트를 구독합니다. 필드명 차이 주의 (`text` → `content`).

***

## 5. 처리 흐름 요약

```
[사용자] ──오디오 업로드──> [Chat API]
                              │
                              ├─ 1. 화자분리 코드서빙 호출
                              │   POST {코드서빙}/json
                              │   body: {audio_url}
                              │   resp: {segments: [{start, end, speaker}]}
                              │
                              ├─ 2. STT 워크플로우 호출 (스트리밍)
                              │   POST {워크플로우}/run/v2
                              │   body: {bucket, object_key, segments[{sequence, start, end}], stream: true}
                              │   stream: {event: utterance, data: {sequence, text}} ... {event: result}
                              │
                              ├─ FE에 utterance_text 이벤트 실시간 전달
                              │
                              ├─ 3. 요약 워크플로우 호출 (스트리밍)
                              │   POST {워크플로우}/run/v2
                              │   body: {question: <transcript>, stream: true}
                              │   stream: {event: token, data: "..."} ... {event: result, data: {text}}
                              │
                              └─ FE에 summary_token / done 이벤트 전달
```

***


---

# Agent Instructions
This documentation is published with GitBook. GitBook is the documentation platform designed so that both humans and AI agents can read, navigate, and reason over technical content effectively. Learn more at gitbook.com.

## Querying This Documentation
If you need additional information that is not directly available in this page, you can query the documentation dynamically by asking a question.

Perform an HTTP GET request on the current page URL with the `ask` query parameter, and the optional `goal` query parameter:

```
GET https://genos-docs.gitbook.io/default/v1.9.2/admin-management/settings/extension-service/ai-meeting-guide.md?ask=<question>&goal=<endgoal>
```

`ask` is the immediate question: it should be specific, self-contained, and written in natural language.
`goal` is optional and describes the broader end goal you are ultimately trying to accomplish on behalf of the user. GitBook uses it to tailor the answer towards what is most useful for that goal.

The response will contain a direct answer to the question and relevant excerpts and sources from the documentation.

Use this mechanism when the answer is not explicitly present in the current page, you need clarification or additional context, or you want to retrieve related documentation sections.
