본문으로 건너뛰기

Wikipedia 편집을 감시하는 로컬 AI 에이전트

Wikipedia 실시간 편집을 값싼 필터와 로컬 LLM으로 감시하는 always-on agent 구축법입니다.

이 요약은 AI가 원문을 분석해 생성했습니다. 정확한 내용은 원문 기준으로 확인하세요.

TL;DR

이 글은 Wikipedia의 실시간 편집 스트림을 감시하면서 로컬 Ollama 모델로 vandalism 후보를 판단하고 그 결과를 토큰 단위로 전송하는 always-on agent 구축법을 다룹니다. 모든 이벤트를 LLM에 보내는 대신 Python 기반 Stage 1에서 삭제된 바이트 수와 사용자별 최근 편집 빈도를 계산해 후보를 줄이고, 통과한 이벤트만 구조화된 AgentVerdict 스키마와 함께 Stage 2 추론에 전달합니다. httpx의 자동 재접속, bounded-memory 기록, 느린 SSE 구독자에 대한 메시지 삭제, FastAPI lifespan 기반 백그라운드 task로 장기 실행 안정성도 확보합니다. 단일 머신을 넘어 여러 소스와 프로세스를 처리하거나 재시작 후 이벤트를 보존하려면 asyncio.Queue를 Kafka 같은 message bus로 바꾸는 확장이 필요합니다.

섹션별 상세

01
이 글은 AI 에이전트에서 혼용되는 두 가지 streaming을 하나의 항상 켜진 시스템으로 결합합니다. 첫 번째는 사람이 메시지를 입력하지 않아도 Wikipedia의 실시간 편집 이벤트를 소비하는 입력 스트리밍이고, 두 번째는 로컬 LLM의 판단 결과를 토큰 단위로 내보내는 출력 스트리밍입니다. 이벤트가 초당 여러 건 발생하는 환경에서 모든 입력을 모델에 보내면 연산 자원이 낭비되고 에이전트가 라이브 피드를 따라가지 못하므로 두 문제를 별도로 처리하는 구조가 필요합니다.
python
def parse_sse_line(line: str) -> Optional[dict]:
    """SSE frames data as lines prefixed with 'data: '. Comment lines (starting with ':') and blank keep-alive lines are common on this feed and should be silently ignored, not treated as errors."""
    if not line or line.startswith(":"):
        return None
    if line.startswith("data:"):
        raw = line[len("data:"):].strip()
        if not raw:
            return None
        try:
            return json.loads(raw)
        except json.JSONDecodeError:
            return None
    return None

SSE 입력에서 주석, 빈 keep-alive 줄과 잘못된 JSON을 걸러내고 데이터 줄만 파싱합니다.

Raw edit events가 Stage 1의 cheap math 필터를 거쳐 Stage 2의 local LLM reasoning으로 전달되는 funnel 구조입니다.
Diagram넓은 영역의 많은 raw edit events 중 대부분은 Stage 1에서 제거되고, 소수의 이벤트만 Stage 2 local LLM reasoning으로 통과합니다. 이후 결과는 Broadcast to live clients 단계로 전달되며, 이는 모든 편집에 모델을 호출하지 않고 삭제량과 편집 빈도 임계값을 통과한 후보만 추론하는 본문의 핵심 설계를 보여줍니다.
02
구현은 인증이 필요 없는 Wikipedia EventStreams를 Server-Sent Events로 읽고, httpx의 비동기 연결이 끊기면 5초 후 재접속하는 방식으로 입력을 지속합니다. 원시 payload에서 edit 타입과 old·new 길이를 확인한 뒤 wiki, 사용자, 문서 제목, 익명 여부, 봇 여부, 편집 요약을 RecentChangeEvent로 정규화하며, 사용자명이 IPv4 또는 IPv6 형태인지 검사해 익명 편집을 판별합니다. parse_sse_line과 to_event를 네트워크와 분리된 순수 함수로 구성해 실제 연결 없이 파싱 로직을 테스트할 수 있게 한 점이 장기 실행 서비스의 복구성과 검증성을 높입니다.
python
class Stage1Filter:
    """Wraps the velocity tracker and the byte-removal check into one pass/fail decision per event."""
    def __init__(self, tracker: Optional[EditVelocityTracker] = None):
        self.tracker = tracker or EditVelocityTracker()

    def evaluate(self, event: RecentChangeEvent) -> Optional[FilterSignal]:
        """Returns a FilterSignal if this event is worth the LLM's time, otherwise None, and None is the common case by a wide margin."""
        if event.is_bot:
            return None # bot edits have their own, separate review path
        recent_count = self.tracker.record_and_count(event.user, event.timestamp)
        bytes_removed = event.bytes_removed
        reasons = []
        if bytes_removed >= config.BYTES_REMOVED_THRESHOLD:
            reasons.append(f"removed {bytes_removed} bytes in one edit")
        if recent_count >= config.EDIT_VELOCITY_THRESHOLD:
            reasons.append(f"{recent_count} edits in {self.tracker.window_seconds}s")
        if not reasons:
            return None
        return FilterSignal(
            event=event,
            bytes_removed=bytes_removed,
            recent_edit_count=recent_count,
            reason="; ".join(reasons),
        )

모든 편집 이벤트에서 봇 여부, 삭제 바이트 수, 최근 편집 빈도를 검사해 LLM 호출 대상을 선별합니다.

03
핵심 병목은 모든 이벤트에 비싼 추론을 적용하는 데 있으므로 2단계 funnel을 사용합니다. Stage 1은 모델 없이 Python으로 삭제된 바이트 수와 sliding window 안의 사용자별 편집 횟수를 계산하고, 봇 편집은 제외한 뒤 임계값을 넘은 이벤트만 FilterSignal로 만듭니다. EditVelocityTracker는 사용자별 deque에서 시간 범위를 벗어난 시각을 제거하고 오래된 사용자 기록을 축출해, 일반 편집 대부분을 값싼 연산으로 버리면서 무한 스트림에서도 메모리가 계속 증가하지 않게 합니다.
python
async def evaluate_signal(signal: FilterSignal) -> AsyncIterator[str | AgentVerdict]:
    """Streams the model's raw output as it's generated (str chunks), then yields a final validated AgentVerdict once the stream completes. The caller tells the two apart with isinstance()."""
    client = ollama.AsyncClient(host=config.OLLAMA_HOST)
    stream = await client.chat(
        model=config.OLLAMA_MODEL,
        messages=[
            {"role": "system", "content": SYSTEM_PROMPT},
            {"role": "user", "content": _build_user_prompt(signal)},
        ],
        format=AgentVerdict.model_json_schema(),
        stream=True,
        options={"temperature": 0.1},
    )
    full_text = ""
    async for chunk in stream:
        piece = chunk["message"]["content"]
        full_text += piece
        if piece:
            yield piece # live token, for the broadcaster to forward immediately
    verdict = AgentVerdict.model_validate_json(full_text)
    yield verdict

Ollama에 JSON 스키마를 전달해 토큰을 실시간으로 내보내고, 완료 후 검증된 AgentVerdict 객체를 반환합니다.

04
Stage 2에서는 Stage 1을 통과한 편집만 Ollama의 로컬 모델로 보내 Wikipedia 문맥에서 vandalism 가능성을 판단합니다. AgentVerdict의 JSON 스키마를 Ollama의 format 인자로 전달하고 temperature를 0.1로 설정해 is_likely_vandalism, 1~5 severity, reasoning, suggested_action을 구조화된 결과로 받으며, 생성 중인 문자열 조각은 즉시 yield하고 완료 후 전체 JSON을 검증합니다. 따라서 연결된 클라이언트는 판단 토큰을 실시간으로 볼 수 있고 후속 코드는 사람이 읽는 임의의 문장이 아니라 타입 검증된 객체를 처리할 수 있습니다.
python
async def publish(self, payload: dict) -> None:
    """Fans a payload out to every subscriber. A subscriber whose queue is full gets the message dropped rather than blocking the whole pipeline, a slow client should never be able to slow down the agent's actual processing loop."""
    message = json.dumps(payload)
    for queue in list(self._subscribers):
        try:
            queue.put_nowait(message)
        except asyncio.QueueFull:
            continue

각 클라이언트의 큐에 메시지를 비동기 전파하며, 느린 구독자의 큐가 가득 차도 전체 처리 루프를 막지 않습니다.

05
Broadcaster는 연결된 클라이언트마다 asyncio.Queue를 만들고 put_nowait으로 flagged, token, verdict 메시지를 독립적으로 전달합니다. 특정 브라우저의 큐가 가득 차면 해당 메시지만 버리고 전체 Wikipedia 소비·추론 루프를 차단하지 않으며, SSE endpoint는 연결 종료를 매번 확인해 구독 큐를 정리합니다. FastAPI lifespan은 애플리케이션 시작 시 단일 백그라운드 pipeline을 만들고 종료 때 취소하므로, 한 번의 Ollama 오류나 느린 클라이언트가 항상 실행되는 서비스를 멈추지 않게 합니다.
python
async def run_pipeline() -> None:
    """Consumes the live stream forever, runs stage 1 on every event, and only calls the LLM stage on events that survive it."""
    async for event in wikipedia_event_stream():
        signal = stage1.evaluate(event)
        if signal is None:
            continue
        logger.info("Stage 1 flagged: %s by %s (%s)", signal.event.title, signal.event.user, signal.reason)
        await broadcaster.publish({"type": "flagged", "title": signal.event.title, "reason": signal.reason})
        try:
            async for item in evaluate_signal(signal):
                if isinstance(item, str):
                    await broadcaster.publish({"type": "token", "title": signal.event.title, "text": item})
                elif isinstance(item, AgentVerdict):
                    await broadcaster.publish({
                        "type": "verdict",
                        "title": signal.event.title,
                        "user": signal.event.user,
                        **item.model_dump(),
                    })
        except Exception:
            logger.exception("Stage 2 failed for %s, skipping this signal", signal.event.title)

실시간 Wikipedia 편집을 1단계 필터에 통과시킨 뒤 후보 이벤트에만 로컬 LLM을 호출하고 결과를 클라이언트에 전송합니다.

06
실행에는 Python 3.11 이상, 로컬 Ollama와 JSON 출력을 지원하는 모델, fastapi·uvicorn·httpx·pydantic·ollama·sse-starlette가 필요하며 클라우드 계정이나 API key는 요구되지 않습니다. 예시 모델은 ollama pull llama3.1:8b로 준비하고 uvicorn src.main:app --reload 뒤 curl -N http://localhost:8000/events로 SSE를 확인합니다. 한 대의 머신에서 하나의 스트림을 감시하는 구성에는 in-process asyncio.Queue와 단일 background task가 맞지만, 다중 소스·다중 프로세스·재시작 후 이벤트 보존이 필요하면 Kafka 같은 message bus로 교체하는 확장 경로가 제시됩니다.

이미지 분석

로컬 장치에서 실행되는 streaming local AI agent의 구성과 데이터 흐름을 나타낸 도식입니다.
Diagram

이미지는 이벤트 입력, 필터·처리 계층, 로컬 컴퓨팅 장치, 클라이언트 출력이 연결된 구조를 시각화합니다. 본문에서 Wikipedia 실시간 편집을 수집하고 Stage 1의 값싼 필터와 Stage 2의 Ollama 기반 로컬 LLM 추론을 거쳐 live clients에 결과를 전달하는 전체 아키텍처와 직접 연결됩니다.

로컬 장치에서 실행되는 streaming local AI agent의 구성과 데이터 흐름을 나타낸 도식입니다.

용어 해설

앰비언트 에이전트(Ambient Agent)
사람이 직접 메시지를 보내야 작동하는 대신 외부 이벤트가 도착하면 자동으로 깨어나는 에이전트입니다. 이벤트 스트림을 감시하다가 조건을 만족한 입력만 처리하므로 요청-응답 방식보다 지속적인 모니터링에 적합합니다. 이 글에서는 Wikipedia 편집 이벤트를 트리거로 사용합니다.
서버 전송 이벤트(Server-Sent Events)
서버가 하나의 HTTP 연결을 유지하면서 클라이언트로 이벤트를 계속 보내는 방식입니다. Wikipedia EventStreams의 편집 데이터를 수신하고, FastAPI 서버가 처리 결과를 브라우저나 curl 클라이언트로 실시간 전달하는 데 사용됩니다. 클라이언트에서 서버로 연결을 유지할 필요가 없는 단방향 스트리밍에 맞습니다.
구조화된 출력(Structured Output)
언어 모델의 응답 형식을 미리 정한 스키마에 맞추는 방식입니다. 이 구현은 AgentVerdict의 JSON 스키마를 Ollama에 전달해 불리언 판정, 1~5 severity, reasoning, suggested_action 필드를 생성하게 합니다. 결과를 사람이 읽는 텍스트가 아니라 프로그램이 바로 처리할 객체로 사용할 수 있습니다.
슬라이딩 윈도(Sliding Window)
최근 일정 시간 동안 발생한 이벤트만 유지해 빈도나 누적량을 계산하는 방식입니다. EditVelocityTracker는 사용자별 편집 시각을 deque에 저장하고 설정된 시간 범위를 벗어난 기록을 제거합니다. 따라서 특정 사용자의 최근 편집 횟수를 계속 갱신하면서 메모리 증가도 제한합니다.
메시지 버스(Message Bus)
이벤트를 생산하는 단계와 소비·추론하는 단계를 분리하는 중간 전달 시스템입니다. 글의 단일 머신 구현은 메모리 기반 asyncio.Queue를 사용하지만, 여러 소스와 프로세스, 재시작 후 미처리 이벤트 보존이 필요할 때 Kafka 같은 메시지 버스로 확장할 수 있습니다. 처리 단계 간 결합을 낮추는 인프라 역할을 합니다.

기술

  • Python 3.11
  • Ollama
  • llama3.1:8b
  • fastapi
  • uvicorn
  • httpx
  • pydantic
  • ollama
  • sse-starlette
  • Wikipedia EventStreams
  • Server-Sent Events
  • FastAPI
  • asyncio
  • Kafka

활용 사례

  • Wikipedia 실시간 편집에서 vandalism 후보 감시
  • 항상 실행되는 이벤트 기반 로컬 AI 에이전트
  • 실시간 모니터링 시스템의 저비용 사전 필터링
  • 브라우저나 curl 클라이언트로 LLM 판단 과정과 verdict 전달
AI 분석 전체 내용 보기

AI 요약 · 북마크 · 개인 피드 설정 — 무료

출처 · 인용 안내

원문 발행 2026. 08. 13.수집 2026. 08. 13.출처 타입 RSS

인용 시 "요약 출처: AI Trends (aitrends.kr)"를 표기하고, 사실 확인은 원문 보기 기준으로 진행해 주세요. 자세한 기준은 운영 정책을 참고해 주세요.