بنقرة واحدة
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 커스텀 이벤트 타입으로 전송 |