一键导入
sse-streaming
Rust Axum SSE 스트리밍 구현 — Sse 응답, Event 구성, tokio-stream 활용, Claude API 스트리밍 변환, EventSource 연동
用 Codex 或 Claude 帮你安装 复制这段 Prompt,粘贴到 Codex、Claude 或其他助手里,让它检查 Skill 页面并帮你完成安装。
菜单
Rust Axum SSE 스트리밍 구현 — Sse 응답, Event 구성, tokio-stream 활용, Claude API 스트리밍 변환, EventSource 연동
用 Codex 或 Claude 帮你安装 复制这段 Prompt,粘贴到 Codex、Claude 或其他助手里,让它检查 Skill 页面并帮你完成安装。
基于 SOC 职业分类
Spring Security 5.5.x + jjwt 0.10.7 레거시 JWT 인증 - WebSecurityConfigurerAdapter, OncePerRequestFilter, javax.servlet 환경
Spring Boot 3.x + Spring Security 6.x + jjwt 0.12.x 기반 모던 JWT 인증 패턴. SecurityFilterChain Bean, 람다 DSL, jakarta.servlet, Virtual Threads 적용
Unity 6 LTS 2D 모바일 게임용 uGUI 시스템 전문 스킬. Canvas/RectTransform/TextMeshPro, 모바일 UI 패턴(팝업·무한 스크롤·광고·IAP), 성능 최적화, UI Toolkit과의 선택 기준 포함.
아크라시아(akrasia, 자제력 없음) 학술 논쟁의 핵심 구도와 주요 연구자·문헌을 빠르게 파악할 수 있는 도메인 지식 스킬. 도덕윤리교육 전공 대학원생(석/박사)이 학위논문·KCI 투고·세미나 준비 시 고대–현대–한국 학계–도덕심리학 흐름을 한 번에 짚도록 구성. <example>사용자: "아리스토텔레스의 propeteia와 astheneia 구분을 인용하려는데 출처를 알려줘"</example> <example>사용자: "데이비슨이 의지박약을 어떻게 가능하다고 봤는지 핵심 논증을 정리해줘"</example> <example>사용자: "한국 도덕교육 학계에서 아크라시아 다룬 논문 있어?"</example>
아리스토텔레스 『니코마코스 윤리학』에서 akrasia(자제력없음)와 akolasia(무절제)의 5축 차이를 정밀하게 정리한 학위논문 자료 스킬. NE VII.4 1147b20-1148b14, VII.8 1150b29-1151a28, III.10-12 1117b23-1119b18 절별 분해와 표준 학자 해석(Bostock, Broadie-Rowe, Pakaluk, Hursthouse, Charles 등)을 포함. 도덕교육 적용을 위한 두 상태 차이의 함의 및 한국어 번역어 처리 권장안 제공. <example>사용자: "akrates와 akolastos를 prohairesis 측면에서 어떻게 구분해야 하나요?"</example> <example>사용자: "NE VII.4의 ἁπλῶς akrasia가 akolasia와 어떻게 갈라지는지 절별 분해해주세요"</example> <example>사용자: "Hursthouse의 연속체 모델을 도덕교육 적용 절에서 어떻게 활용할 수 있나요?"</example>
한국 위기 대응 자원(자살·자해·정신건강·여성·청소년·노인·다문화) 핫라인과 앱·챗봇 안전 가드 응답 패턴 종합. 꿈 해몽·정신건강 앱 등 자가 진단/감정 콘텐츠 도메인에서 위험 신호 포착 시 안전한 자원 안내 문구를 작성할 때 참조. <example>사용자: "꿈 해몽 앱에 위기 안내 문구를 어떻게 넣을까?"</example> <example>사용자: "한국에서 자살예방 핫라인 번호가 어떻게 바뀌었지?"</example> <example>사용자: "정신건강 챗봇 안전 가드 응답 템플릿을 짜줘"</example>
| name | sse-streaming |
| description | Rust Axum SSE 스트리밍 구현 — Sse 응답, Event 구성, tokio-stream 활용, Claude API 스트리밍 변환, EventSource 연동 |
소스: https://docs.rs/axum/latest/axum/response/sse/index.html | https://docs.rs/tokio-stream/latest/tokio_stream/ | https://developer.mozilla.org/en-US/docs/Web/API/EventSource 검증일: 2026-06-20
주의: Axum 0.8.x 기준으로 작성. 0.7 이하에서는
axum::response::sse모듈 경로 및 일부 API가 다를 수 있다.
# Cargo.toml
[dependencies]
axum = "0.8"
tokio = { version = "1", features = ["full"] }
tokio-stream = "0.1"
futures = "0.3"
serde = { version = "1", features = ["derive"] }
serde_json = "1"
axum::response::sse::Sse는 IntoResponse를 구현하는 SSE 응답 래퍼다. Stream<Item = Result<Event, E>>를 받아 SSE 프로토콜로 변환한다.
use axum::response::sse::{Event, Sse};
use futures::stream::Stream;
use std::convert::Infallible;
async fn sse_handler() -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
let stream = tokio_stream::iter(vec![
Ok(Event::default().data("hello")),
Ok(Event::default().data("world")),
]);
Sse::new(stream)
}
핵심: Sse::new(stream)에 전달하는 스트림의 Item은 Result<Event, E>여야 한다. E가 에러 없는 스트림이면 Infallible을 사용한다.
Event는 SSE 프로토콜의 단일 이벤트를 나타낸다. 빌더 패턴으로 구성한다.
use axum::response::sse::Event;
// 기본 데이터 전송
Event::default().data("plain text message")
// 이벤트 타입 지정 — 프론트엔드에서 addEventListener로 수신
Event::default()
.event("chat_message")
.data("hello from server")
// id 설정 — 클라이언트 재연결 시 Last-Event-ID 헤더로 전송됨
Event::default()
.id("msg-001")
.data("tracked message")
// JSON 데이터 전송
let payload = serde_json::json!({ "role": "assistant", "content": "Hi" });
Event::default()
.event("message")
.data(payload.to_string())
// 재연결 간격 설정 (밀리초)
Event::default()
.retry(std::time::Duration::from_secs(5))
.data("retry configured")
// 주석 — keep-alive용으로 활용
Event::default().comment("keep-alive")
| 메서드 | SSE 필드 | 용도 |
|---|---|---|
.data(str) | data: | 이벤트 페이로드 |
.event(str) | event: | 이벤트 타입명 (기본: message) |
.id(str) | id: | 이벤트 ID (재연결 추적) |
.retry(Duration) | retry: | 재연결 대기 시간 |
.comment(str) | : | 주석 (keep-alive 등) |
비동기 작업에서 이벤트를 push할 때 가장 일반적인 패턴이다.
use axum::response::sse::{Event, Sse};
use tokio_stream::wrappers::ReceiverStream;
use tokio::sync::mpsc;
use std::convert::Infallible;
async fn sse_handler() -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
let (tx, rx) = mpsc::channel::<Result<Event, Infallible>>(100);
tokio::spawn(async move {
for i in 0..10 {
let event = Event::default()
.event("counter")
.data(format!("{}", i));
if tx.send(Ok(event)).await.is_err() {
break; // 클라이언트 연결 종료
}
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
});
Sse::new(ReceiverStream::new(rx))
}
간결한 스트림 생성이 필요할 때 사용한다.
use async_stream::stream;
use axum::response::sse::{Event, Sse};
use std::convert::Infallible;
async fn sse_handler() -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
let s = stream! {
for i in 0..10 {
yield Ok(Event::default().data(format!("count: {}", i)));
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
};
Sse::new(s)
}
주의:
async-stream은 별도 크레이트(async-stream = "0.3")가 필요하다. Axum 공식 예제는stream!매크로 대신futures_util::stream::repeat_with()를 사용한다.stream!은 커뮤니티에서 널리 사용되는 패턴이지만 공식 예제 포함 여부와는 무관하다.
연결 유지를 위해 주기적으로 주석 이벤트를 보낸다.
use axum::response::sse::KeepAlive;
Sse::new(stream).keep_alive(
KeepAlive::new()
.interval(std::time::Duration::from_secs(15))
.text("keep-alive")
)
use axum::{routing::get, Router};
let app = Router::new()
.route("/events", get(sse_handler));
SSE는 GET 요청으로 처리한다. POST 본문이 필요한 경우(예: Claude API 호출 파라미터 전달) 별도 엔드포인트로 세션을 생성하고, SSE 연결에서 세션 ID를 쿼리로 전달하는 패턴을 사용한다.
Claude Messages API의 stream: true 응답을 서버에서 수신하여 프론트엔드로 중계하는 패턴이다.
use axum::{
extract::Json,
response::sse::{Event, Sse},
};
use futures::stream::Stream;
use reqwest::Client;
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use std::convert::Infallible;
#[derive(serde::Deserialize)]
struct ChatRequest {
message: String,
}
async fn chat_stream(
Json(req): Json<ChatRequest>,
) -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
let (tx, rx) = mpsc::channel::<Result<Event, Infallible>>(100);
tokio::spawn(async move {
let client = Client::new();
let api_key = std::env::var("ANTHROPIC_API_KEY")
.expect("ANTHROPIC_API_KEY not set");
// Claude Messages API 스트리밍 요청
let response = client
.post("https://api.anthropic.com/v1/messages")
.header("x-api-key", &api_key)
.header("anthropic-version", "2023-06-01")
.header("content-type", "application/json")
.json(&serde_json::json!({
"model": "claude-sonnet-4-6",
"max_tokens": 1024,
"stream": true,
"messages": [{"role": "user", "content": req.message}]
}))
.send()
.await;
let Ok(resp) = response else {
let _ = tx.send(Ok(
Event::default().event("error").data("API request failed")
)).await;
return;
};
// 바이트 스트림을 SSE 라인으로 파싱
let mut stream = resp.bytes_stream();
let mut buffer = String::new();
use futures::StreamExt;
while let Some(chunk) = stream.next().await {
let Ok(bytes) = chunk else { break };
buffer.push_str(&String::from_utf8_lossy(&bytes));
// Claude SSE 응답은 "data: {...}\n\n" 형식
while let Some(pos) = buffer.find("\n\n") {
let line = buffer[..pos].to_string();
buffer = buffer[pos + 2..].to_string();
let Some(data) = line.strip_prefix("data: ") else {
continue;
};
// event: message_stop 이면 종료
if data.contains("\"type\":\"message_stop\"") {
let _ = tx.send(Ok(
Event::default().event("done").data("[DONE]")
)).await;
return;
}
// content_block_delta에서 텍스트 추출
if data.contains("\"type\":\"content_block_delta\"") {
if let Ok(parsed) = serde_json::from_str::<serde_json::Value>(data) {
if let Some(text) = parsed["delta"]["text"].as_str() {
let _ = tx.send(Ok(
Event::default().event("text_delta").data(text)
)).await;
}
}
}
}
}
});
Sse::new(ReceiverStream::new(rx))
.keep_alive(
axum::response::sse::KeepAlive::new()
.interval(std::time::Duration::from_secs(15))
)
}
주의: Claude API 스트리밍 응답의 전체 이벤트 타입은
message_start,content_block_start,content_block_delta,content_block_stop,message_delta,message_stop,ping,error이며, extended thinking 사용 시thinking_delta,signature_delta가 추가됩니다. API 버전에 따라 변경될 수 있으므로 최신 사양은 https://docs.anthropic.com/en/api/messages-streaming 참조.
브라우저 EventSource가 다른 오리진의 SSE에 연결하려면 CORS 설정이 필요하다.
use tower_http::cors::{CorsLayer, Any};
let cors = CorsLayer::new()
.allow_origin(Any)
.allow_methods(Any)
.allow_headers(Any);
let app = Router::new()
.route("/events", get(sse_handler))
.layer(cors);
주의: 프로덕션에서는
Any대신 허용할 오리진을 명시적으로 지정한다.
const source = new EventSource("/events");
// 기본 message 이벤트
source.onmessage = (e: MessageEvent) => {
console.log("data:", e.data);
};
// 커스텀 이벤트 타입 수신 (.event()로 지정한 이름)
source.addEventListener("text_delta", (e: MessageEvent) => {
console.log("delta:", e.data);
});
source.addEventListener("done", () => {
source.close();
});
// 에러 핸들링
source.onerror = (e) => {
console.error("SSE error:", e);
source.close();
};
EventSource는 GET만 지원한다. POST 본문이 필요하면 fetch의 ReadableStream을 파싱한다.
async function streamChat(message: string, onDelta: (text: string) => void) {
const response = await fetch("/api/chat", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ message }),
});
if (!response.ok || !response.body) {
throw new Error(`HTTP ${response.status}`);
}
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
// SSE 프로토콜 파싱: "event: ...\ndata: ...\n\n"
const parts = buffer.split("\n\n");
buffer = parts.pop() ?? "";
for (const part of parts) {
const lines = part.split("\n");
let eventType = "message";
let data = "";
for (const line of lines) {
if (line.startsWith("event: ")) eventType = line.slice(7);
if (line.startsWith("data: ")) data = line.slice(6);
}
if (eventType === "text_delta") {
onDelta(data);
} else if (eventType === "done") {
return;
}
}
}
}
// 사용 예
streamChat("Hello Claude", (text) => {
process.stdout.write(text); // 또는 DOM에 append
});
import { useState, useCallback } from "react";
function useSSEChat() {
const [content, setContent] = useState("");
const [isStreaming, setIsStreaming] = useState(false);
const send = useCallback(async (message: string) => {
setContent("");
setIsStreaming(true);
await streamChat(message, (delta) => {
setContent((prev) => prev + delta);
});
setIsStreaming(false);
}, []);
return { content, isStreaming, send };
}
use axum::response::sse::Event;
// 에러를 SSE 이벤트로 변환
fn error_event(msg: &str) -> Event {
Event::default()
.event("error")
.data(serde_json::json!({ "message": msg }).to_string())
}
// 스트림 내에서 에러 발생 시
if let Err(e) = some_operation().await {
let _ = tx.send(Ok(error_event(&e.to_string()))).await;
return;
}
| 항목 | 선택 |
|---|---|
| 응답 타입 | Sse<impl Stream<Item = Result<Event, Infallible>>> |
| 스트림 생성 | mpsc + ReceiverStream (비동기 push) 또는 async_stream::stream! (간결) |
| Keep-Alive | Sse::new(s).keep_alive(KeepAlive::new().interval(...)) |
| 프론트 GET | EventSource API |
| 프론트 POST | fetch + ReadableStream 파싱 |
| 에러 전달 | event: error 커스텀 이벤트 타입으로 전송 |