update: Merge worker-types crate into worker crate
This commit is contained in:
+1
-2
@@ -13,8 +13,7 @@ serde_json = "1.0"
|
||||
thiserror = "1.0"
|
||||
tokio = { version = "1.49.0", features = ["macros", "rt-multi-thread"] }
|
||||
tracing = "0.1"
|
||||
worker-macros = { path = "../worker-macros" }
|
||||
worker-types = { path = "../worker-types" }
|
||||
worker-macros = { path = "../worker-macros", version = "0.1" }
|
||||
|
||||
[dev-dependencies]
|
||||
clap = { version = "4.5.54", features = ["derive", "env"] }
|
||||
|
||||
@@ -0,0 +1,446 @@
|
||||
//! Worker層の公開イベント型
|
||||
//!
|
||||
//! 外部利用者に公開するためのイベント表現。
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
// =============================================================================
|
||||
// Core Event Types (from llm_client layer)
|
||||
// =============================================================================
|
||||
|
||||
/// LLMからのストリーミングイベント
|
||||
///
|
||||
/// 各LLMプロバイダからのレスポンスは、この`Event`のストリームとして
|
||||
/// 統一的に処理されます。
|
||||
///
|
||||
/// # イベントの種類
|
||||
///
|
||||
/// - **メタイベント**: `Ping`, `Usage`, `Status`, `Error`
|
||||
/// - **ブロックイベント**: `BlockStart`, `BlockDelta`, `BlockStop`, `BlockAbort`
|
||||
///
|
||||
/// # ブロックのライフサイクル
|
||||
///
|
||||
/// テキストやツール呼び出しは、`BlockStart` → `BlockDelta`(複数) → `BlockStop`
|
||||
/// の順序でイベントが発生します。
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum Event {
|
||||
/// ハートビート
|
||||
Ping(PingEvent),
|
||||
/// トークン使用量
|
||||
Usage(UsageEvent),
|
||||
/// ストリームのステータス変化
|
||||
Status(StatusEvent),
|
||||
/// エラー発生
|
||||
Error(ErrorEvent),
|
||||
|
||||
/// ブロック開始(テキスト、ツール使用等)
|
||||
BlockStart(BlockStart),
|
||||
/// ブロックの差分データ
|
||||
BlockDelta(BlockDelta),
|
||||
/// ブロック正常終了
|
||||
BlockStop(BlockStop),
|
||||
/// ブロック中断
|
||||
BlockAbort(BlockAbort),
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Meta Events
|
||||
// =============================================================================
|
||||
|
||||
/// Pingイベント(ハートビート)
|
||||
#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)]
|
||||
pub struct PingEvent {
|
||||
pub timestamp: Option<u64>,
|
||||
}
|
||||
|
||||
/// 使用量イベント
|
||||
#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)]
|
||||
pub struct UsageEvent {
|
||||
/// 入力トークン数
|
||||
pub input_tokens: Option<u64>,
|
||||
/// 出力トークン数
|
||||
pub output_tokens: Option<u64>,
|
||||
/// 合計トークン数
|
||||
pub total_tokens: Option<u64>,
|
||||
/// キャッシュ読み込みトークン数
|
||||
pub cache_read_input_tokens: Option<u64>,
|
||||
/// キャッシュ作成トークン数
|
||||
pub cache_creation_input_tokens: Option<u64>,
|
||||
}
|
||||
|
||||
/// ステータスイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct StatusEvent {
|
||||
pub status: ResponseStatus,
|
||||
}
|
||||
|
||||
/// レスポンスステータス
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum ResponseStatus {
|
||||
/// ストリーム開始
|
||||
Started,
|
||||
/// 正常完了
|
||||
Completed,
|
||||
/// キャンセルされた
|
||||
Cancelled,
|
||||
/// エラー発生
|
||||
Failed,
|
||||
}
|
||||
|
||||
/// エラーイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct ErrorEvent {
|
||||
pub code: Option<String>,
|
||||
pub message: String,
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Block Types
|
||||
// =============================================================================
|
||||
|
||||
/// ブロックの種別
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
pub enum BlockType {
|
||||
/// テキスト生成
|
||||
Text,
|
||||
/// 思考 (Claude Extended Thinking等)
|
||||
Thinking,
|
||||
/// ツール呼び出し
|
||||
ToolUse,
|
||||
/// ツール結果
|
||||
ToolResult,
|
||||
}
|
||||
|
||||
/// ブロック開始イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockStart {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// ブロック固有のメタデータ
|
||||
pub metadata: BlockMetadata,
|
||||
}
|
||||
|
||||
impl BlockStart {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロックのメタデータ
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum BlockMetadata {
|
||||
Text,
|
||||
Thinking,
|
||||
ToolUse { id: String, name: String },
|
||||
ToolResult { tool_use_id: String },
|
||||
}
|
||||
|
||||
/// ブロックデルタイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockDelta {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// デルタの内容
|
||||
pub delta: DeltaContent,
|
||||
}
|
||||
|
||||
/// デルタの内容
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum DeltaContent {
|
||||
/// テキストデルタ
|
||||
Text(String),
|
||||
/// 思考デルタ
|
||||
Thinking(String),
|
||||
/// ツール引数のJSON部分文字列
|
||||
InputJson(String),
|
||||
}
|
||||
|
||||
impl DeltaContent {
|
||||
/// デルタのブロック種別を取得
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
match self {
|
||||
DeltaContent::Text(_) => BlockType::Text,
|
||||
DeltaContent::Thinking(_) => BlockType::Thinking,
|
||||
DeltaContent::InputJson(_) => BlockType::ToolUse,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロック停止イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockStop {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// 停止理由
|
||||
pub stop_reason: Option<StopReason>,
|
||||
}
|
||||
|
||||
impl BlockStop {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロック中断イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockAbort {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// 中断理由
|
||||
pub reason: String,
|
||||
}
|
||||
|
||||
impl BlockAbort {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// 停止理由
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum StopReason {
|
||||
/// 自然終了
|
||||
EndTurn,
|
||||
/// 最大トークン数到達
|
||||
MaxTokens,
|
||||
/// ストップシーケンス到達
|
||||
StopSequence,
|
||||
/// ツール使用
|
||||
ToolUse,
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Builder / Factory helpers
|
||||
// =============================================================================
|
||||
|
||||
impl Event {
|
||||
/// テキストブロック開始イベントを作成
|
||||
pub fn text_block_start(index: usize) -> Self {
|
||||
Event::BlockStart(BlockStart {
|
||||
index,
|
||||
block_type: BlockType::Text,
|
||||
metadata: BlockMetadata::Text,
|
||||
})
|
||||
}
|
||||
|
||||
/// テキストデルタイベントを作成
|
||||
pub fn text_delta(index: usize, text: impl Into<String>) -> Self {
|
||||
Event::BlockDelta(BlockDelta {
|
||||
index,
|
||||
delta: DeltaContent::Text(text.into()),
|
||||
})
|
||||
}
|
||||
|
||||
/// テキストブロック停止イベントを作成
|
||||
pub fn text_block_stop(index: usize, stop_reason: Option<StopReason>) -> Self {
|
||||
Event::BlockStop(BlockStop {
|
||||
index,
|
||||
block_type: BlockType::Text,
|
||||
stop_reason,
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール使用ブロック開始イベントを作成
|
||||
pub fn tool_use_start(index: usize, id: impl Into<String>, name: impl Into<String>) -> Self {
|
||||
Event::BlockStart(BlockStart {
|
||||
index,
|
||||
block_type: BlockType::ToolUse,
|
||||
metadata: BlockMetadata::ToolUse {
|
||||
id: id.into(),
|
||||
name: name.into(),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール引数デルタイベントを作成
|
||||
pub fn tool_input_delta(index: usize, json: impl Into<String>) -> Self {
|
||||
Event::BlockDelta(BlockDelta {
|
||||
index,
|
||||
delta: DeltaContent::InputJson(json.into()),
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール使用ブロック停止イベントを作成
|
||||
pub fn tool_use_stop(index: usize) -> Self {
|
||||
Event::BlockStop(BlockStop {
|
||||
index,
|
||||
block_type: BlockType::ToolUse,
|
||||
stop_reason: Some(StopReason::ToolUse),
|
||||
})
|
||||
}
|
||||
|
||||
/// 使用量イベントを作成
|
||||
pub fn usage(input_tokens: u64, output_tokens: u64) -> Self {
|
||||
Event::Usage(UsageEvent {
|
||||
input_tokens: Some(input_tokens),
|
||||
output_tokens: Some(output_tokens),
|
||||
total_tokens: Some(input_tokens + output_tokens),
|
||||
cache_read_input_tokens: None,
|
||||
cache_creation_input_tokens: None,
|
||||
})
|
||||
}
|
||||
|
||||
/// Pingイベントを作成
|
||||
pub fn ping() -> Self {
|
||||
Event::Ping(PingEvent { timestamp: None })
|
||||
}
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Conversions: timeline::event -> worker::event
|
||||
// =============================================================================
|
||||
|
||||
impl From<crate::timeline::event::ResponseStatus> for ResponseStatus {
|
||||
fn from(value: crate::timeline::event::ResponseStatus) -> Self {
|
||||
match value {
|
||||
crate::timeline::event::ResponseStatus::Started => ResponseStatus::Started,
|
||||
crate::timeline::event::ResponseStatus::Completed => ResponseStatus::Completed,
|
||||
crate::timeline::event::ResponseStatus::Cancelled => ResponseStatus::Cancelled,
|
||||
crate::timeline::event::ResponseStatus::Failed => ResponseStatus::Failed,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::BlockType> for BlockType {
|
||||
fn from(value: crate::timeline::event::BlockType) -> Self {
|
||||
match value {
|
||||
crate::timeline::event::BlockType::Text => BlockType::Text,
|
||||
crate::timeline::event::BlockType::Thinking => BlockType::Thinking,
|
||||
crate::timeline::event::BlockType::ToolUse => BlockType::ToolUse,
|
||||
crate::timeline::event::BlockType::ToolResult => BlockType::ToolResult,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::BlockMetadata> for BlockMetadata {
|
||||
fn from(value: crate::timeline::event::BlockMetadata) -> Self {
|
||||
match value {
|
||||
crate::timeline::event::BlockMetadata::Text => BlockMetadata::Text,
|
||||
crate::timeline::event::BlockMetadata::Thinking => BlockMetadata::Thinking,
|
||||
crate::timeline::event::BlockMetadata::ToolUse { id, name } => {
|
||||
BlockMetadata::ToolUse { id, name }
|
||||
}
|
||||
crate::timeline::event::BlockMetadata::ToolResult { tool_use_id } => {
|
||||
BlockMetadata::ToolResult { tool_use_id }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::DeltaContent> for DeltaContent {
|
||||
fn from(value: crate::timeline::event::DeltaContent) -> Self {
|
||||
match value {
|
||||
crate::timeline::event::DeltaContent::Text(text) => DeltaContent::Text(text),
|
||||
crate::timeline::event::DeltaContent::Thinking(text) => DeltaContent::Thinking(text),
|
||||
crate::timeline::event::DeltaContent::InputJson(json) => DeltaContent::InputJson(json),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::StopReason> for StopReason {
|
||||
fn from(value: crate::timeline::event::StopReason) -> Self {
|
||||
match value {
|
||||
crate::timeline::event::StopReason::EndTurn => StopReason::EndTurn,
|
||||
crate::timeline::event::StopReason::MaxTokens => StopReason::MaxTokens,
|
||||
crate::timeline::event::StopReason::StopSequence => StopReason::StopSequence,
|
||||
crate::timeline::event::StopReason::ToolUse => StopReason::ToolUse,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::PingEvent> for PingEvent {
|
||||
fn from(value: crate::timeline::event::PingEvent) -> Self {
|
||||
PingEvent {
|
||||
timestamp: value.timestamp,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::UsageEvent> for UsageEvent {
|
||||
fn from(value: crate::timeline::event::UsageEvent) -> Self {
|
||||
UsageEvent {
|
||||
input_tokens: value.input_tokens,
|
||||
output_tokens: value.output_tokens,
|
||||
total_tokens: value.total_tokens,
|
||||
cache_read_input_tokens: value.cache_read_input_tokens,
|
||||
cache_creation_input_tokens: value.cache_creation_input_tokens,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::StatusEvent> for StatusEvent {
|
||||
fn from(value: crate::timeline::event::StatusEvent) -> Self {
|
||||
StatusEvent {
|
||||
status: value.status.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::ErrorEvent> for ErrorEvent {
|
||||
fn from(value: crate::timeline::event::ErrorEvent) -> Self {
|
||||
ErrorEvent {
|
||||
code: value.code,
|
||||
message: value.message,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::BlockStart> for BlockStart {
|
||||
fn from(value: crate::timeline::event::BlockStart) -> Self {
|
||||
BlockStart {
|
||||
index: value.index,
|
||||
block_type: value.block_type.into(),
|
||||
metadata: value.metadata.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::BlockDelta> for BlockDelta {
|
||||
fn from(value: crate::timeline::event::BlockDelta) -> Self {
|
||||
BlockDelta {
|
||||
index: value.index,
|
||||
delta: value.delta.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::BlockStop> for BlockStop {
|
||||
fn from(value: crate::timeline::event::BlockStop) -> Self {
|
||||
BlockStop {
|
||||
index: value.index,
|
||||
block_type: value.block_type.into(),
|
||||
stop_reason: value.stop_reason.map(Into::into),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::BlockAbort> for BlockAbort {
|
||||
fn from(value: crate::timeline::event::BlockAbort) -> Self {
|
||||
BlockAbort {
|
||||
index: value.index,
|
||||
block_type: value.block_type.into(),
|
||||
reason: value.reason,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::timeline::event::Event> for Event {
|
||||
fn from(value: crate::timeline::event::Event) -> Self {
|
||||
match value {
|
||||
crate::timeline::event::Event::Ping(p) => Event::Ping(p.into()),
|
||||
crate::timeline::event::Event::Usage(u) => Event::Usage(u.into()),
|
||||
crate::timeline::event::Event::Status(s) => Event::Status(s.into()),
|
||||
crate::timeline::event::Event::Error(e) => Event::Error(e.into()),
|
||||
crate::timeline::event::Event::BlockStart(s) => Event::BlockStart(s.into()),
|
||||
crate::timeline::event::Event::BlockDelta(d) => Event::BlockDelta(d.into()),
|
||||
crate::timeline::event::Event::BlockStop(s) => Event::BlockStop(s.into()),
|
||||
crate::timeline::event::Event::BlockAbort(a) => Event::BlockAbort(a.into()),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,174 @@
|
||||
//! Handler/Kind型
|
||||
//!
|
||||
//! Timeline層でイベントを処理するためのトレイト。
|
||||
//! カスタムハンドラを実装してTimelineに登録することで、
|
||||
//! ストリームイベントを受信できます。
|
||||
|
||||
use crate::timeline::event::*;
|
||||
|
||||
// =============================================================================
|
||||
// Kind Trait
|
||||
// =============================================================================
|
||||
|
||||
/// イベント種別を定義するマーカートレイト
|
||||
///
|
||||
/// 各Kindは対応するイベント型を指定します。
|
||||
/// HandlerはこのKindに対して実装され、同じKindに対して
|
||||
/// 異なるScope型を持つ複数のHandlerを登録できます。
|
||||
pub trait Kind {
|
||||
/// このKindに対応するイベント型
|
||||
type Event;
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Handler Trait
|
||||
// =============================================================================
|
||||
|
||||
/// イベントを処理するハンドラトレイト
|
||||
///
|
||||
/// 特定の`Kind`に対するイベント処理を定義します。
|
||||
/// `Scope`はブロックのライフサイクル中に保持される状態です。
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::timeline::{Handler, TextBlockEvent, TextBlockKind};
|
||||
///
|
||||
/// struct TextCollector {
|
||||
/// texts: Vec<String>,
|
||||
/// }
|
||||
///
|
||||
/// impl Handler<TextBlockKind> for TextCollector {
|
||||
/// type Scope = String; // ブロックごとのバッファ
|
||||
///
|
||||
/// fn on_event(&mut self, buffer: &mut String, event: &TextBlockEvent) {
|
||||
/// match event {
|
||||
/// TextBlockEvent::Delta(text) => buffer.push_str(text),
|
||||
/// TextBlockEvent::Stop(_) => {
|
||||
/// self.texts.push(std::mem::take(buffer));
|
||||
/// }
|
||||
/// _ => {}
|
||||
/// }
|
||||
/// }
|
||||
/// }
|
||||
/// ```
|
||||
pub trait Handler<K: Kind> {
|
||||
/// Handler固有のスコープ型
|
||||
///
|
||||
/// ブロック開始時に`Default::default()`で生成され、
|
||||
/// ブロック終了時に破棄されます。
|
||||
type Scope: Default;
|
||||
|
||||
/// イベントを処理する
|
||||
fn on_event(&mut self, scope: &mut Self::Scope, event: &K::Event);
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Meta Kind Definitions
|
||||
// =============================================================================
|
||||
|
||||
/// Usage Kind - 使用量イベント用
|
||||
pub struct UsageKind;
|
||||
impl Kind for UsageKind {
|
||||
type Event = UsageEvent;
|
||||
}
|
||||
|
||||
/// Ping Kind - Pingイベント用
|
||||
pub struct PingKind;
|
||||
impl Kind for PingKind {
|
||||
type Event = PingEvent;
|
||||
}
|
||||
|
||||
/// Status Kind - ステータスイベント用
|
||||
pub struct StatusKind;
|
||||
impl Kind for StatusKind {
|
||||
type Event = StatusEvent;
|
||||
}
|
||||
|
||||
/// Error Kind - エラーイベント用
|
||||
pub struct ErrorKind;
|
||||
impl Kind for ErrorKind {
|
||||
type Event = ErrorEvent;
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Block Kind Definitions
|
||||
// =============================================================================
|
||||
|
||||
/// TextBlock Kind - テキストブロック用
|
||||
pub struct TextBlockKind;
|
||||
impl Kind for TextBlockKind {
|
||||
type Event = TextBlockEvent;
|
||||
}
|
||||
|
||||
/// テキストブロックのイベント
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub enum TextBlockEvent {
|
||||
Start(TextBlockStart),
|
||||
Delta(String),
|
||||
Stop(TextBlockStop),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct TextBlockStart {
|
||||
pub index: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct TextBlockStop {
|
||||
pub index: usize,
|
||||
pub stop_reason: Option<StopReason>,
|
||||
}
|
||||
|
||||
/// ThinkingBlock Kind - 思考ブロック用
|
||||
pub struct ThinkingBlockKind;
|
||||
impl Kind for ThinkingBlockKind {
|
||||
type Event = ThinkingBlockEvent;
|
||||
}
|
||||
|
||||
/// 思考ブロックのイベント
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub enum ThinkingBlockEvent {
|
||||
Start(ThinkingBlockStart),
|
||||
Delta(String),
|
||||
Stop(ThinkingBlockStop),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct ThinkingBlockStart {
|
||||
pub index: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct ThinkingBlockStop {
|
||||
pub index: usize,
|
||||
}
|
||||
|
||||
/// ToolUseBlock Kind - ツール使用ブロック用
|
||||
pub struct ToolUseBlockKind;
|
||||
impl Kind for ToolUseBlockKind {
|
||||
type Event = ToolUseBlockEvent;
|
||||
}
|
||||
|
||||
/// ツール使用ブロックのイベント
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub enum ToolUseBlockEvent {
|
||||
Start(ToolUseBlockStart),
|
||||
/// ツール引数のJSON部分文字列
|
||||
InputJsonDelta(String),
|
||||
Stop(ToolUseBlockStop),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct ToolUseBlockStart {
|
||||
pub index: usize,
|
||||
pub id: String,
|
||||
pub name: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct ToolUseBlockStop {
|
||||
pub index: usize,
|
||||
pub id: String,
|
||||
pub name: String,
|
||||
}
|
||||
@@ -0,0 +1,182 @@
|
||||
//! Hook関連の型定義
|
||||
//!
|
||||
//! Worker層でのターン制御・介入に使用される型
|
||||
|
||||
use async_trait::async_trait;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde_json::Value;
|
||||
use thiserror::Error;
|
||||
|
||||
// =============================================================================
|
||||
// Control Flow Types
|
||||
// =============================================================================
|
||||
|
||||
/// Hook処理の制御フロー
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum ControlFlow {
|
||||
/// 処理を続行
|
||||
Continue,
|
||||
/// 現在の処理をスキップ(Tool実行など)
|
||||
Skip,
|
||||
/// 処理を中断
|
||||
Abort(String),
|
||||
}
|
||||
|
||||
/// ターン終了時の判定結果
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum TurnResult {
|
||||
/// ターンを終了
|
||||
Finish,
|
||||
/// メッセージを追加してターン継続(自己修正など)
|
||||
ContinueWithMessages(Vec<crate::Message>),
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Tool Call / Result Types
|
||||
// =============================================================================
|
||||
|
||||
/// ツール呼び出し情報
|
||||
///
|
||||
/// LLMからのToolUseブロックを表現し、Hook処理で改変可能
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ToolCall {
|
||||
/// ツール呼び出しID(レスポンスとの紐付けに使用)
|
||||
pub id: String,
|
||||
/// ツール名
|
||||
pub name: String,
|
||||
/// 入力引数(JSON)
|
||||
pub input: Value,
|
||||
}
|
||||
|
||||
/// ツール実行結果
|
||||
///
|
||||
/// ツール実行後の結果を表現し、Hook処理で改変可能
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ToolResult {
|
||||
/// 対応するツール呼び出しID
|
||||
pub tool_use_id: String,
|
||||
/// 結果コンテンツ
|
||||
pub content: String,
|
||||
/// エラーかどうか
|
||||
#[serde(default)]
|
||||
pub is_error: bool,
|
||||
}
|
||||
|
||||
impl ToolResult {
|
||||
/// 成功結果を作成
|
||||
pub fn success(tool_use_id: impl Into<String>, content: impl Into<String>) -> Self {
|
||||
Self {
|
||||
tool_use_id: tool_use_id.into(),
|
||||
content: content.into(),
|
||||
is_error: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// エラー結果を作成
|
||||
pub fn error(tool_use_id: impl Into<String>, content: impl Into<String>) -> Self {
|
||||
Self {
|
||||
tool_use_id: tool_use_id.into(),
|
||||
content: content.into(),
|
||||
is_error: true,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Hook Error
|
||||
// =============================================================================
|
||||
|
||||
/// Hookエラー
|
||||
#[derive(Debug, Error)]
|
||||
pub enum HookError {
|
||||
/// 処理が中断された
|
||||
#[error("Aborted: {0}")]
|
||||
Aborted(String),
|
||||
/// 内部エラー
|
||||
#[error("Hook error: {0}")]
|
||||
Internal(String),
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// WorkerHook Trait
|
||||
// =============================================================================
|
||||
|
||||
/// ターンの進行・ツール実行に介入するためのトレイト
|
||||
///
|
||||
/// Hookを使うと、メッセージ送信前、ツール実行前後、ターン終了時に
|
||||
/// 処理を挟んだり、実行をキャンセルしたりできます。
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::hook::{ControlFlow, HookError, ToolCall, TurnResult, WorkerHook};
|
||||
/// use worker::Message;
|
||||
///
|
||||
/// struct ValidationHook;
|
||||
///
|
||||
/// #[async_trait::async_trait]
|
||||
/// impl WorkerHook for ValidationHook {
|
||||
/// async fn before_tool_call(&self, call: &mut ToolCall) -> Result<ControlFlow, HookError> {
|
||||
/// // 危険なツールをブロック
|
||||
/// if call.name == "delete_all" {
|
||||
/// return Ok(ControlFlow::Skip);
|
||||
/// }
|
||||
/// Ok(ControlFlow::Continue)
|
||||
/// }
|
||||
///
|
||||
/// async fn on_turn_end(&self, messages: &[Message]) -> Result<TurnResult, HookError> {
|
||||
/// // 条件を満たさなければ追加メッセージで継続
|
||||
/// if messages.len() < 3 {
|
||||
/// return Ok(TurnResult::ContinueWithMessages(vec![
|
||||
/// Message::user("Please elaborate.")
|
||||
/// ]));
|
||||
/// }
|
||||
/// Ok(TurnResult::Finish)
|
||||
/// }
|
||||
/// }
|
||||
/// ```
|
||||
///
|
||||
/// # デフォルト実装
|
||||
///
|
||||
/// すべてのメソッドにはデフォルト実装があり、何も行わず`Continue`を返します。
|
||||
/// 必要なメソッドのみオーバーライドしてください。
|
||||
#[async_trait]
|
||||
pub trait WorkerHook: Send + Sync {
|
||||
/// メッセージ送信前に呼ばれる
|
||||
///
|
||||
/// リクエストに含まれるメッセージリストを参照・改変できます。
|
||||
/// `ControlFlow::Abort`を返すとターンが中断されます。
|
||||
async fn on_message_send(
|
||||
&self,
|
||||
_context: &mut Vec<crate::Message>,
|
||||
) -> Result<ControlFlow, HookError> {
|
||||
Ok(ControlFlow::Continue)
|
||||
}
|
||||
|
||||
/// ツール実行前に呼ばれる
|
||||
///
|
||||
/// ツール呼び出しの引数を書き換えたり、実行をスキップしたりできます。
|
||||
/// `ControlFlow::Skip`を返すとこのツールの実行がスキップされます。
|
||||
async fn before_tool_call(&self, _tool_call: &mut ToolCall) -> Result<ControlFlow, HookError> {
|
||||
Ok(ControlFlow::Continue)
|
||||
}
|
||||
|
||||
/// ツール実行後に呼ばれる
|
||||
///
|
||||
/// ツールの実行結果を書き換えたり、隠蔽したりできます。
|
||||
async fn after_tool_call(
|
||||
&self,
|
||||
_tool_result: &mut ToolResult,
|
||||
) -> Result<ControlFlow, HookError> {
|
||||
Ok(ControlFlow::Continue)
|
||||
}
|
||||
|
||||
/// ターン終了時に呼ばれる
|
||||
///
|
||||
/// 生成されたメッセージを検査し、必要なら追加メッセージで継続を指示できます。
|
||||
/// `TurnResult::ContinueWithMessages`を返すと、指定したメッセージを追加して
|
||||
/// 次のターンに進みます。
|
||||
async fn on_turn_end(&self, _messages: &[crate::Message]) -> Result<TurnResult, HookError> {
|
||||
Ok(TurnResult::Finish)
|
||||
}
|
||||
}
|
||||
+8
-45
@@ -39,55 +39,18 @@
|
||||
pub mod llm_client;
|
||||
pub mod timeline;
|
||||
|
||||
mod subscriber_adapter;
|
||||
pub mod event;
|
||||
mod handler;
|
||||
pub mod hook;
|
||||
mod message;
|
||||
pub mod state;
|
||||
pub mod subscriber;
|
||||
pub mod tool;
|
||||
mod worker;
|
||||
|
||||
// =============================================================================
|
||||
// トップレベル公開(最も頻繁に使う型)
|
||||
// =============================================================================
|
||||
|
||||
pub use message::{ContentPart, Message, MessageContent, Role};
|
||||
pub use worker::{Worker, WorkerConfig, WorkerError};
|
||||
pub use worker_types::{ContentPart, Message, MessageContent, Role};
|
||||
|
||||
// =============================================================================
|
||||
// 意味のあるモジュールとして公開
|
||||
// =============================================================================
|
||||
|
||||
/// ツール定義
|
||||
///
|
||||
/// LLMから呼び出し可能なツールを定義するためのトレイトと型。
|
||||
pub mod tool {
|
||||
pub use worker_types::{Tool, ToolError};
|
||||
}
|
||||
|
||||
/// Hook機能
|
||||
///
|
||||
/// ターンの進行・ツール実行に介入するためのトレイトと型。
|
||||
pub mod hook {
|
||||
pub use worker_types::{ControlFlow, HookError, ToolCall, ToolResult, TurnResult, WorkerHook};
|
||||
}
|
||||
|
||||
/// イベント購読
|
||||
///
|
||||
/// LLMからのストリーミングイベントをリアルタイムで受信するためのトレイト。
|
||||
pub mod subscriber {
|
||||
pub use worker_types::WorkerSubscriber;
|
||||
}
|
||||
|
||||
/// イベント型
|
||||
///
|
||||
/// LLMからのストリーミングレスポンスを表現するイベント型。
|
||||
/// Timeline層を直接使用する場合に必要です。
|
||||
pub mod event {
|
||||
pub use worker_types::{
|
||||
BlockAbort, BlockDelta, BlockMetadata, BlockStart, BlockStop, BlockType, DeltaContent,
|
||||
ErrorEvent, Event, PingEvent, ResponseStatus, StatusEvent, StopReason, UsageEvent,
|
||||
};
|
||||
}
|
||||
|
||||
/// Worker状態
|
||||
///
|
||||
/// Type-stateパターンによるキャッシュ保護のための状態マーカー型。
|
||||
pub mod state {
|
||||
pub use worker_types::{Locked, Mutable, WorkerState};
|
||||
}
|
||||
|
||||
@@ -2,11 +2,9 @@
|
||||
|
||||
use std::pin::Pin;
|
||||
|
||||
use crate::llm_client::{ClientError, Request, event::Event};
|
||||
use async_trait::async_trait;
|
||||
use futures::Stream;
|
||||
use worker_types::Event;
|
||||
|
||||
use crate::llm_client::{ClientError, Request};
|
||||
|
||||
/// LLMクライアントのtrait
|
||||
///
|
||||
|
||||
@@ -0,0 +1,293 @@
|
||||
//! LLMクライアント層のイベント型
|
||||
//!
|
||||
//! 各LLMプロバイダからのストリーミングレスポンスを表現するイベント型。
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
// =============================================================================
|
||||
// Core Event Types (from llm_client layer)
|
||||
// =============================================================================
|
||||
|
||||
/// LLMからのストリーミングイベント
|
||||
///
|
||||
/// 各LLMプロバイダからのレスポンスは、この`Event`のストリームとして
|
||||
/// 統一的に処理されます。
|
||||
///
|
||||
/// # イベントの種類
|
||||
///
|
||||
/// - **メタイベント**: `Ping`, `Usage`, `Status`, `Error`
|
||||
/// - **ブロックイベント**: `BlockStart`, `BlockDelta`, `BlockStop`, `BlockAbort`
|
||||
///
|
||||
/// # ブロックのライフサイクル
|
||||
///
|
||||
/// テキストやツール呼び出しは、`BlockStart` → `BlockDelta`(複数) → `BlockStop`
|
||||
/// の順序でイベントが発生します。
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum Event {
|
||||
/// ハートビート
|
||||
Ping(PingEvent),
|
||||
/// トークン使用量
|
||||
Usage(UsageEvent),
|
||||
/// ストリームのステータス変化
|
||||
Status(StatusEvent),
|
||||
/// エラー発生
|
||||
Error(ErrorEvent),
|
||||
|
||||
/// ブロック開始(テキスト、ツール使用等)
|
||||
BlockStart(BlockStart),
|
||||
/// ブロックの差分データ
|
||||
BlockDelta(BlockDelta),
|
||||
/// ブロック正常終了
|
||||
BlockStop(BlockStop),
|
||||
/// ブロック中断
|
||||
BlockAbort(BlockAbort),
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Meta Events
|
||||
// =============================================================================
|
||||
|
||||
/// Pingイベント(ハートビート)
|
||||
#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)]
|
||||
pub struct PingEvent {
|
||||
pub timestamp: Option<u64>,
|
||||
}
|
||||
|
||||
/// 使用量イベント
|
||||
#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)]
|
||||
pub struct UsageEvent {
|
||||
/// 入力トークン数
|
||||
pub input_tokens: Option<u64>,
|
||||
/// 出力トークン数
|
||||
pub output_tokens: Option<u64>,
|
||||
/// 合計トークン数
|
||||
pub total_tokens: Option<u64>,
|
||||
/// キャッシュ読み込みトークン数
|
||||
pub cache_read_input_tokens: Option<u64>,
|
||||
/// キャッシュ作成トークン数
|
||||
pub cache_creation_input_tokens: Option<u64>,
|
||||
}
|
||||
|
||||
/// ステータスイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct StatusEvent {
|
||||
pub status: ResponseStatus,
|
||||
}
|
||||
|
||||
/// レスポンスステータス
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum ResponseStatus {
|
||||
/// ストリーム開始
|
||||
Started,
|
||||
/// 正常完了
|
||||
Completed,
|
||||
/// キャンセルされた
|
||||
Cancelled,
|
||||
/// エラー発生
|
||||
Failed,
|
||||
}
|
||||
|
||||
/// エラーイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct ErrorEvent {
|
||||
pub code: Option<String>,
|
||||
pub message: String,
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Block Types
|
||||
// =============================================================================
|
||||
|
||||
/// ブロックの種別
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
pub enum BlockType {
|
||||
/// テキスト生成
|
||||
Text,
|
||||
/// 思考 (Claude Extended Thinking等)
|
||||
Thinking,
|
||||
/// ツール呼び出し
|
||||
ToolUse,
|
||||
/// ツール結果
|
||||
ToolResult,
|
||||
}
|
||||
|
||||
/// ブロック開始イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockStart {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// ブロック固有のメタデータ
|
||||
pub metadata: BlockMetadata,
|
||||
}
|
||||
|
||||
impl BlockStart {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロックのメタデータ
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum BlockMetadata {
|
||||
Text,
|
||||
Thinking,
|
||||
ToolUse { id: String, name: String },
|
||||
ToolResult { tool_use_id: String },
|
||||
}
|
||||
|
||||
/// ブロックデルタイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockDelta {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// デルタの内容
|
||||
pub delta: DeltaContent,
|
||||
}
|
||||
|
||||
/// デルタの内容
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum DeltaContent {
|
||||
/// テキストデルタ
|
||||
Text(String),
|
||||
/// 思考デルタ
|
||||
Thinking(String),
|
||||
/// ツール引数のJSON部分文字列
|
||||
InputJson(String),
|
||||
}
|
||||
|
||||
impl DeltaContent {
|
||||
/// デルタのブロック種別を取得
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
match self {
|
||||
DeltaContent::Text(_) => BlockType::Text,
|
||||
DeltaContent::Thinking(_) => BlockType::Thinking,
|
||||
DeltaContent::InputJson(_) => BlockType::ToolUse,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロック停止イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockStop {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// 停止理由
|
||||
pub stop_reason: Option<StopReason>,
|
||||
}
|
||||
|
||||
impl BlockStop {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロック中断イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockAbort {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// 中断理由
|
||||
pub reason: String,
|
||||
}
|
||||
|
||||
impl BlockAbort {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// 停止理由
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum StopReason {
|
||||
/// 自然終了
|
||||
EndTurn,
|
||||
/// 最大トークン数到達
|
||||
MaxTokens,
|
||||
/// ストップシーケンス到達
|
||||
StopSequence,
|
||||
/// ツール使用
|
||||
ToolUse,
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Builder / Factory helpers
|
||||
// =============================================================================
|
||||
|
||||
impl Event {
|
||||
/// テキストブロック開始イベントを作成
|
||||
pub fn text_block_start(index: usize) -> Self {
|
||||
Event::BlockStart(BlockStart {
|
||||
index,
|
||||
block_type: BlockType::Text,
|
||||
metadata: BlockMetadata::Text,
|
||||
})
|
||||
}
|
||||
|
||||
/// テキストデルタイベントを作成
|
||||
pub fn text_delta(index: usize, text: impl Into<String>) -> Self {
|
||||
Event::BlockDelta(BlockDelta {
|
||||
index,
|
||||
delta: DeltaContent::Text(text.into()),
|
||||
})
|
||||
}
|
||||
|
||||
/// テキストブロック停止イベントを作成
|
||||
pub fn text_block_stop(index: usize, stop_reason: Option<StopReason>) -> Self {
|
||||
Event::BlockStop(BlockStop {
|
||||
index,
|
||||
block_type: BlockType::Text,
|
||||
stop_reason,
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール使用ブロック開始イベントを作成
|
||||
pub fn tool_use_start(index: usize, id: impl Into<String>, name: impl Into<String>) -> Self {
|
||||
Event::BlockStart(BlockStart {
|
||||
index,
|
||||
block_type: BlockType::ToolUse,
|
||||
metadata: BlockMetadata::ToolUse {
|
||||
id: id.into(),
|
||||
name: name.into(),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール引数デルタイベントを作成
|
||||
pub fn tool_input_delta(index: usize, json: impl Into<String>) -> Self {
|
||||
Event::BlockDelta(BlockDelta {
|
||||
index,
|
||||
delta: DeltaContent::InputJson(json.into()),
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール使用ブロック停止イベントを作成
|
||||
pub fn tool_use_stop(index: usize) -> Self {
|
||||
Event::BlockStop(BlockStop {
|
||||
index,
|
||||
block_type: BlockType::ToolUse,
|
||||
stop_reason: Some(StopReason::ToolUse),
|
||||
})
|
||||
}
|
||||
|
||||
/// 使用量イベントを作成
|
||||
pub fn usage(input_tokens: u64, output_tokens: u64) -> Self {
|
||||
Event::Usage(UsageEvent {
|
||||
input_tokens: Some(input_tokens),
|
||||
output_tokens: Some(output_tokens),
|
||||
total_tokens: Some(input_tokens + output_tokens),
|
||||
cache_read_input_tokens: None,
|
||||
cache_creation_input_tokens: None,
|
||||
})
|
||||
}
|
||||
|
||||
/// Pingイベントを作成
|
||||
pub fn ping() -> Self {
|
||||
Event::Ping(PingEvent { timestamp: None })
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
//! LLMクライアント層
|
||||
//!
|
||||
//! 各LLMプロバイダと通信し、統一された[`Event`](crate::event::Event)ストリームを出力します。
|
||||
//! 各LLMプロバイダと通信し、統一された[`Event`](crate::llm_client::event::Event)
|
||||
//! ストリームを出力します。
|
||||
//!
|
||||
//! # サポートするプロバイダ
|
||||
//!
|
||||
@@ -17,6 +18,7 @@
|
||||
|
||||
pub mod client;
|
||||
pub mod error;
|
||||
pub mod event;
|
||||
pub mod types;
|
||||
|
||||
pub mod providers;
|
||||
@@ -24,4 +26,5 @@ pub mod scheme;
|
||||
|
||||
pub use client::*;
|
||||
pub use error::*;
|
||||
pub use event::*;
|
||||
pub use types::*;
|
||||
|
||||
@@ -4,13 +4,13 @@
|
||||
|
||||
use std::pin::Pin;
|
||||
|
||||
use crate::llm_client::{
|
||||
ClientError, LlmClient, Request, event::Event, scheme::anthropic::AnthropicScheme,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use eventsource_stream::Eventsource;
|
||||
use futures::{Stream, StreamExt, TryStreamExt, future::ready};
|
||||
use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderValue};
|
||||
use worker_types::Event;
|
||||
|
||||
use crate::llm_client::{ClientError, LlmClient, Request, scheme::anthropic::AnthropicScheme};
|
||||
|
||||
/// Anthropic クライアント
|
||||
pub struct AnthropicClient {
|
||||
@@ -156,10 +156,11 @@ impl LlmClient for AnthropicClient {
|
||||
if let Some(block_type) = current_block_type.take() {
|
||||
// 正しいブロックタイプで上書き
|
||||
// (Event::BlockStopの中身を置換)
|
||||
evt = Event::BlockStop(worker_types::BlockStop {
|
||||
block_type,
|
||||
..stop.clone()
|
||||
});
|
||||
evt =
|
||||
Event::BlockStop(crate::llm_client::event::BlockStop {
|
||||
block_type,
|
||||
..stop.clone()
|
||||
});
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
|
||||
@@ -4,13 +4,13 @@
|
||||
|
||||
use std::pin::Pin;
|
||||
|
||||
use crate::llm_client::{
|
||||
ClientError, LlmClient, Request, event::Event, scheme::gemini::GeminiScheme,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use eventsource_stream::Eventsource;
|
||||
use futures::{Stream, StreamExt, TryStreamExt};
|
||||
use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderValue};
|
||||
use worker_types::Event;
|
||||
|
||||
use crate::llm_client::{ClientError, LlmClient, Request, scheme::gemini::GeminiScheme};
|
||||
|
||||
/// Gemini クライアント
|
||||
pub struct GeminiClient {
|
||||
|
||||
@@ -5,13 +5,12 @@
|
||||
|
||||
use std::pin::Pin;
|
||||
|
||||
use crate::llm_client::{
|
||||
ClientError, LlmClient, Request, event::Event, providers::openai::OpenAIClient,
|
||||
scheme::openai::OpenAIScheme,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use futures::Stream;
|
||||
use worker_types::Event;
|
||||
|
||||
use crate::llm_client::{
|
||||
ClientError, LlmClient, Request, providers::openai::OpenAIClient, scheme::openai::OpenAIScheme,
|
||||
};
|
||||
|
||||
/// Ollama クライアント
|
||||
///
|
||||
|
||||
@@ -4,13 +4,13 @@
|
||||
|
||||
use std::pin::Pin;
|
||||
|
||||
use crate::llm_client::{
|
||||
ClientError, LlmClient, Request, event::Event, scheme::openai::OpenAIScheme,
|
||||
};
|
||||
use async_trait::async_trait;
|
||||
use eventsource_stream::Eventsource;
|
||||
use futures::{Stream, StreamExt, TryStreamExt};
|
||||
use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderValue};
|
||||
use worker_types::Event;
|
||||
|
||||
use crate::llm_client::{ClientError, LlmClient, Request, scheme::openai::OpenAIScheme};
|
||||
|
||||
/// OpenAI クライアント
|
||||
pub struct OpenAIClient {
|
||||
|
||||
@@ -2,13 +2,14 @@
|
||||
//!
|
||||
//! Anthropic Messages APIのSSEイベントをパースし、統一Event型に変換
|
||||
|
||||
use serde::Deserialize;
|
||||
use worker_types::{
|
||||
BlockDelta, BlockMetadata, BlockStart, BlockStop, BlockType, DeltaContent, ErrorEvent, Event,
|
||||
PingEvent, ResponseStatus, StatusEvent, UsageEvent,
|
||||
use crate::llm_client::{
|
||||
ClientError,
|
||||
event::{
|
||||
BlockDelta, BlockMetadata, BlockStart, BlockStop, BlockType, DeltaContent, ErrorEvent,
|
||||
Event, PingEvent, ResponseStatus, StatusEvent, UsageEvent,
|
||||
},
|
||||
};
|
||||
|
||||
use crate::llm_client::ClientError;
|
||||
use serde::Deserialize;
|
||||
|
||||
use super::AnthropicScheme;
|
||||
|
||||
|
||||
@@ -2,12 +2,11 @@
|
||||
//!
|
||||
//! Google Gemini APIのSSEイベントをパースし、統一Event型に変換
|
||||
|
||||
use serde::Deserialize;
|
||||
use worker_types::{
|
||||
BlockMetadata, BlockStart, BlockStop, BlockType, Event, StopReason, UsageEvent,
|
||||
use crate::llm_client::{
|
||||
ClientError,
|
||||
event::{BlockMetadata, BlockStart, BlockStop, BlockType, Event, StopReason, UsageEvent},
|
||||
};
|
||||
|
||||
use crate::llm_client::ClientError;
|
||||
use serde::Deserialize;
|
||||
|
||||
use super::GeminiScheme;
|
||||
|
||||
@@ -231,7 +230,7 @@ impl GeminiScheme {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use worker_types::DeltaContent;
|
||||
use crate::llm_client::event::DeltaContent;
|
||||
|
||||
#[test]
|
||||
fn test_parse_text_response() {
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
//! OpenAI SSEイベントパース
|
||||
|
||||
use crate::llm_client::{
|
||||
ClientError,
|
||||
event::{Event, StopReason, UsageEvent},
|
||||
};
|
||||
use serde::Deserialize;
|
||||
use worker_types::{Event, StopReason, UsageEvent};
|
||||
|
||||
use crate::llm_client::ClientError;
|
||||
|
||||
use super::OpenAIScheme;
|
||||
|
||||
@@ -155,7 +156,7 @@ impl OpenAIScheme {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use worker_types::DeltaContent;
|
||||
use crate::llm_client::event::DeltaContent;
|
||||
|
||||
#[test]
|
||||
fn test_parse_text_delta() {
|
||||
@@ -188,7 +189,7 @@ mod tests {
|
||||
assert_eq!(events.len(), 1);
|
||||
if let Event::BlockStart(start) = &events[0] {
|
||||
assert_eq!(start.index, 0);
|
||||
if let worker_types::BlockMetadata::ToolUse { id, name } = &start.metadata {
|
||||
if let crate::llm_client::event::BlockMetadata::ToolUse { id, name } = &start.metadata {
|
||||
assert_eq!(id, "call_abc");
|
||||
assert_eq!(name, "get_weather");
|
||||
} else {
|
||||
|
||||
@@ -0,0 +1,116 @@
|
||||
//! メッセージ型
|
||||
//!
|
||||
//! LLMとの会話で使用されるメッセージ構造。
|
||||
//! [`Message::user`]や[`Message::assistant`]で簡単に作成できます。
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
/// メッセージのロール
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
pub enum Role {
|
||||
/// ユーザー
|
||||
User,
|
||||
/// アシスタント
|
||||
Assistant,
|
||||
}
|
||||
|
||||
/// 会話のメッセージ
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::Message;
|
||||
///
|
||||
/// // ユーザーメッセージ
|
||||
/// let user_msg = Message::user("Hello!");
|
||||
///
|
||||
/// // アシスタントメッセージ
|
||||
/// let assistant_msg = Message::assistant("Hi there!");
|
||||
/// ```
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct Message {
|
||||
/// ロール
|
||||
pub role: Role,
|
||||
/// コンテンツ
|
||||
pub content: MessageContent,
|
||||
}
|
||||
|
||||
/// メッセージコンテンツ
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(untagged)]
|
||||
pub enum MessageContent {
|
||||
/// テキストコンテンツ
|
||||
Text(String),
|
||||
/// ツール結果
|
||||
ToolResult {
|
||||
tool_use_id: String,
|
||||
content: String,
|
||||
},
|
||||
/// 複合コンテンツ (テキスト + ツール使用等)
|
||||
Parts(Vec<ContentPart>),
|
||||
}
|
||||
|
||||
/// コンテンツパーツ
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
#[serde(tag = "type")]
|
||||
pub enum ContentPart {
|
||||
/// テキスト
|
||||
#[serde(rename = "text")]
|
||||
Text { text: String },
|
||||
/// ツール使用
|
||||
#[serde(rename = "tool_use")]
|
||||
ToolUse {
|
||||
id: String,
|
||||
name: String,
|
||||
input: serde_json::Value,
|
||||
},
|
||||
/// ツール結果
|
||||
#[serde(rename = "tool_result")]
|
||||
ToolResult {
|
||||
tool_use_id: String,
|
||||
content: String,
|
||||
},
|
||||
}
|
||||
|
||||
impl Message {
|
||||
/// ユーザーメッセージを作成
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::Message;
|
||||
/// let msg = Message::user("こんにちは");
|
||||
/// ```
|
||||
pub fn user(content: impl Into<String>) -> Self {
|
||||
Self {
|
||||
role: Role::User,
|
||||
content: MessageContent::Text(content.into()),
|
||||
}
|
||||
}
|
||||
|
||||
/// アシスタントメッセージを作成
|
||||
///
|
||||
/// 通常はWorker内部で自動生成されますが、
|
||||
/// 履歴の初期化などで手動作成も可能です。
|
||||
pub fn assistant(content: impl Into<String>) -> Self {
|
||||
Self {
|
||||
role: Role::Assistant,
|
||||
content: MessageContent::Text(content.into()),
|
||||
}
|
||||
}
|
||||
|
||||
/// ツール結果メッセージを作成
|
||||
///
|
||||
/// Worker内部でツール実行後に自動生成されます。
|
||||
/// 通常は直接作成する必要はありません。
|
||||
pub fn tool_result(tool_use_id: impl Into<String>, content: impl Into<String>) -> Self {
|
||||
Self {
|
||||
role: Role::User,
|
||||
content: MessageContent::ToolResult {
|
||||
tool_use_id: tool_use_id.into(),
|
||||
content: content.into(),
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
//! Worker状態
|
||||
//!
|
||||
//! Type-stateパターンによるキャッシュ保護のための状態マーカー型。
|
||||
//! Workerは`Mutable` → `Locked`の状態遷移を持ちます。
|
||||
|
||||
/// Worker状態を表すマーカートレイト
|
||||
///
|
||||
/// このトレイトはシールされており、外部から実装することはできません。
|
||||
pub trait WorkerState: private::Sealed + Send + Sync + 'static {}
|
||||
|
||||
mod private {
|
||||
pub trait Sealed {}
|
||||
}
|
||||
|
||||
/// 編集可能状態
|
||||
///
|
||||
/// この状態では以下の操作が可能です:
|
||||
/// - システムプロンプトの設定・変更
|
||||
/// - メッセージ履歴の編集(追加、削除、クリア)
|
||||
/// - ツール・Hookの登録
|
||||
///
|
||||
/// `Worker::lock()`により[`Locked`]状態へ遷移できます。
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::Worker;
|
||||
///
|
||||
/// let mut worker = Worker::new(client)
|
||||
/// .system_prompt("You are helpful.");
|
||||
///
|
||||
/// // 履歴を編集可能
|
||||
/// worker.push_message(Message::user("Hello"));
|
||||
/// worker.clear_history();
|
||||
///
|
||||
/// // ロックして保護状態へ
|
||||
/// let locked = worker.lock();
|
||||
/// ```
|
||||
#[derive(Debug, Clone, Copy, Default)]
|
||||
pub struct Mutable;
|
||||
|
||||
impl private::Sealed for Mutable {}
|
||||
impl WorkerState for Mutable {}
|
||||
|
||||
/// ロック状態(キャッシュ保護)
|
||||
///
|
||||
/// この状態では以下の制限があります:
|
||||
/// - システムプロンプトの変更不可
|
||||
/// - 既存メッセージ履歴の変更不可(末尾への追記のみ)
|
||||
///
|
||||
/// LLM APIのKVキャッシュヒットを保証するため、
|
||||
/// 実行時にはこの状態の使用が推奨されます。
|
||||
///
|
||||
/// `Worker::unlock()`により[`Mutable`]状態へ戻せますが、
|
||||
/// キャッシュ保護が解除されることに注意してください。
|
||||
#[derive(Debug, Clone, Copy, Default)]
|
||||
pub struct Locked;
|
||||
|
||||
impl private::Sealed for Locked {}
|
||||
impl WorkerState for Locked {}
|
||||
@@ -1,14 +1,145 @@
|
||||
//! WorkerSubscriber統合
|
||||
//! イベント購読
|
||||
//!
|
||||
//! WorkerSubscriberをTimeline層のHandlerとしてブリッジする実装
|
||||
//! LLMからのストリーミングイベントをリアルタイムで受信するためのトレイト。
|
||||
//! UIへのストリーム表示やプログレス表示に使用します。
|
||||
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use worker_types::{
|
||||
ErrorEvent, ErrorKind, Handler, StatusEvent, StatusKind, TextBlockEvent, TextBlockKind,
|
||||
ToolCall, ToolUseBlockEvent, ToolUseBlockKind, UsageEvent, UsageKind, WorkerSubscriber,
|
||||
use crate::{
|
||||
handler::{
|
||||
ErrorKind, Handler, StatusKind, TextBlockEvent, TextBlockKind, ToolUseBlockEvent,
|
||||
ToolUseBlockKind, UsageKind,
|
||||
},
|
||||
hook::ToolCall,
|
||||
timeline::event::{ErrorEvent, StatusEvent, UsageEvent},
|
||||
};
|
||||
|
||||
// =============================================================================
|
||||
// WorkerSubscriber Trait
|
||||
// =============================================================================
|
||||
|
||||
/// LLMからのストリーミングイベントを購読するトレイト
|
||||
///
|
||||
/// Workerに登録すると、テキスト生成やツール呼び出しのイベントを
|
||||
/// リアルタイムで受信できます。UIへのストリーム表示に最適です。
|
||||
///
|
||||
/// # 受信できるイベント
|
||||
///
|
||||
/// - **ブロックイベント**: テキスト、ツール使用(スコープ付き)
|
||||
/// - **メタイベント**: 使用量、ステータス、エラー
|
||||
/// - **完了イベント**: テキスト完了、ツール呼び出し完了
|
||||
/// - **ターン制御**: ターン開始、ターン終了
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::subscriber::WorkerSubscriber;
|
||||
/// use worker::timeline::TextBlockEvent;
|
||||
///
|
||||
/// struct StreamPrinter;
|
||||
///
|
||||
/// impl WorkerSubscriber for StreamPrinter {
|
||||
/// type TextBlockScope = ();
|
||||
/// type ToolUseBlockScope = ();
|
||||
///
|
||||
/// fn on_text_block(&mut self, _: &mut (), event: &TextBlockEvent) {
|
||||
/// if let TextBlockEvent::Delta(text) = event {
|
||||
/// print!("{}", text); // リアルタイム出力
|
||||
/// }
|
||||
/// }
|
||||
///
|
||||
/// fn on_text_complete(&mut self, text: &str) {
|
||||
/// println!("\n--- Complete: {} chars ---", text.len());
|
||||
/// }
|
||||
/// }
|
||||
///
|
||||
/// // Workerに登録
|
||||
/// worker.subscribe(StreamPrinter);
|
||||
/// ```
|
||||
pub trait WorkerSubscriber: Send {
|
||||
// =========================================================================
|
||||
// スコープ型(ブロックイベント用)
|
||||
// =========================================================================
|
||||
|
||||
/// テキストブロック処理用のスコープ型
|
||||
///
|
||||
/// ブロック開始時にDefault::default()で生成され、
|
||||
/// ブロック終了時に破棄される。
|
||||
type TextBlockScope: Default + Send;
|
||||
|
||||
/// ツール使用ブロック処理用のスコープ型
|
||||
type ToolUseBlockScope: Default + Send;
|
||||
|
||||
// =========================================================================
|
||||
// ブロックイベント(スコープ管理あり)
|
||||
// =========================================================================
|
||||
|
||||
/// テキストブロックイベント
|
||||
///
|
||||
/// Start/Delta/Stopのライフサイクルを持つ。
|
||||
/// scopeはブロック開始時に生成され、終了時に破棄される。
|
||||
#[allow(unused_variables)]
|
||||
fn on_text_block(&mut self, scope: &mut Self::TextBlockScope, event: &TextBlockEvent) {}
|
||||
|
||||
/// ツール使用ブロックイベント
|
||||
///
|
||||
/// Start/InputJsonDelta/Stopのライフサイクルを持つ。
|
||||
#[allow(unused_variables)]
|
||||
fn on_tool_use_block(
|
||||
&mut self,
|
||||
scope: &mut Self::ToolUseBlockScope,
|
||||
event: &ToolUseBlockEvent,
|
||||
) {
|
||||
}
|
||||
|
||||
// =========================================================================
|
||||
// 単発イベント(スコープ不要)
|
||||
// =========================================================================
|
||||
|
||||
/// 使用量イベント
|
||||
#[allow(unused_variables)]
|
||||
fn on_usage(&mut self, event: &UsageEvent) {}
|
||||
|
||||
/// ステータスイベント
|
||||
#[allow(unused_variables)]
|
||||
fn on_status(&mut self, event: &StatusEvent) {}
|
||||
|
||||
/// エラーイベント
|
||||
#[allow(unused_variables)]
|
||||
fn on_error(&mut self, event: &ErrorEvent) {}
|
||||
|
||||
// =========================================================================
|
||||
// 累積イベント(Worker層で追加)
|
||||
// =========================================================================
|
||||
|
||||
/// テキスト完了イベント
|
||||
///
|
||||
/// テキストブロックが完了した時点で、累積されたテキスト全体が渡される。
|
||||
/// ブロック処理後の最終結果を受け取るのに便利。
|
||||
#[allow(unused_variables)]
|
||||
fn on_text_complete(&mut self, text: &str) {}
|
||||
|
||||
/// ツール呼び出し完了イベント
|
||||
///
|
||||
/// ツール使用ブロックが完了した時点で、完全なToolCallが渡される。
|
||||
#[allow(unused_variables)]
|
||||
fn on_tool_call_complete(&mut self, call: &ToolCall) {}
|
||||
|
||||
// =========================================================================
|
||||
// ターン制御
|
||||
// =========================================================================
|
||||
|
||||
/// ターン開始時
|
||||
///
|
||||
/// `turn`は0から始まるターン番号。
|
||||
#[allow(unused_variables)]
|
||||
fn on_turn_start(&mut self, turn: usize) {}
|
||||
|
||||
/// ターン終了時
|
||||
#[allow(unused_variables)]
|
||||
fn on_turn_end(&mut self, turn: usize) {}
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// SubscriberAdapter - WorkerSubscriberをTimelineハンドラにブリッジ
|
||||
// =============================================================================
|
||||
@@ -0,0 +1,448 @@
|
||||
//! Timeline層のイベント型
|
||||
//!
|
||||
//! Timelineが受け取り、各Handlerへディスパッチするイベント表現。
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
// =============================================================================
|
||||
// Core Event Types (from llm_client layer)
|
||||
// =============================================================================
|
||||
|
||||
/// LLMからのストリーミングイベント
|
||||
///
|
||||
/// 各LLMプロバイダからのレスポンスは、この`Event`のストリームとして
|
||||
/// 統一的に処理されます。
|
||||
///
|
||||
/// # イベントの種類
|
||||
///
|
||||
/// - **メタイベント**: `Ping`, `Usage`, `Status`, `Error`
|
||||
/// - **ブロックイベント**: `BlockStart`, `BlockDelta`, `BlockStop`, `BlockAbort`
|
||||
///
|
||||
/// # ブロックのライフサイクル
|
||||
///
|
||||
/// テキストやツール呼び出しは、`BlockStart` → `BlockDelta`(複数) → `BlockStop`
|
||||
/// の順序でイベントが発生します。
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum Event {
|
||||
/// ハートビート
|
||||
Ping(PingEvent),
|
||||
/// トークン使用量
|
||||
Usage(UsageEvent),
|
||||
/// ストリームのステータス変化
|
||||
Status(StatusEvent),
|
||||
/// エラー発生
|
||||
Error(ErrorEvent),
|
||||
|
||||
/// ブロック開始(テキスト、ツール使用等)
|
||||
BlockStart(BlockStart),
|
||||
/// ブロックの差分データ
|
||||
BlockDelta(BlockDelta),
|
||||
/// ブロック正常終了
|
||||
BlockStop(BlockStop),
|
||||
/// ブロック中断
|
||||
BlockAbort(BlockAbort),
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Meta Events
|
||||
// =============================================================================
|
||||
|
||||
/// Pingイベント(ハートビート)
|
||||
#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)]
|
||||
pub struct PingEvent {
|
||||
pub timestamp: Option<u64>,
|
||||
}
|
||||
|
||||
/// 使用量イベント
|
||||
#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)]
|
||||
pub struct UsageEvent {
|
||||
/// 入力トークン数
|
||||
pub input_tokens: Option<u64>,
|
||||
/// 出力トークン数
|
||||
pub output_tokens: Option<u64>,
|
||||
/// 合計トークン数
|
||||
pub total_tokens: Option<u64>,
|
||||
/// キャッシュ読み込みトークン数
|
||||
pub cache_read_input_tokens: Option<u64>,
|
||||
/// キャッシュ作成トークン数
|
||||
pub cache_creation_input_tokens: Option<u64>,
|
||||
}
|
||||
|
||||
/// ステータスイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct StatusEvent {
|
||||
pub status: ResponseStatus,
|
||||
}
|
||||
|
||||
/// レスポンスステータス
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum ResponseStatus {
|
||||
/// ストリーム開始
|
||||
Started,
|
||||
/// 正常完了
|
||||
Completed,
|
||||
/// キャンセルされた
|
||||
Cancelled,
|
||||
/// エラー発生
|
||||
Failed,
|
||||
}
|
||||
|
||||
/// エラーイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct ErrorEvent {
|
||||
pub code: Option<String>,
|
||||
pub message: String,
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Block Types
|
||||
// =============================================================================
|
||||
|
||||
/// ブロックの種別
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
pub enum BlockType {
|
||||
/// テキスト生成
|
||||
Text,
|
||||
/// 思考 (Claude Extended Thinking等)
|
||||
Thinking,
|
||||
/// ツール呼び出し
|
||||
ToolUse,
|
||||
/// ツール結果
|
||||
ToolResult,
|
||||
}
|
||||
|
||||
/// ブロック開始イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockStart {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// ブロック固有のメタデータ
|
||||
pub metadata: BlockMetadata,
|
||||
}
|
||||
|
||||
impl BlockStart {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロックのメタデータ
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum BlockMetadata {
|
||||
Text,
|
||||
Thinking,
|
||||
ToolUse { id: String, name: String },
|
||||
ToolResult { tool_use_id: String },
|
||||
}
|
||||
|
||||
/// ブロックデルタイベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockDelta {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// デルタの内容
|
||||
pub delta: DeltaContent,
|
||||
}
|
||||
|
||||
/// デルタの内容
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum DeltaContent {
|
||||
/// テキストデルタ
|
||||
Text(String),
|
||||
/// 思考デルタ
|
||||
Thinking(String),
|
||||
/// ツール引数のJSON部分文字列
|
||||
InputJson(String),
|
||||
}
|
||||
|
||||
impl DeltaContent {
|
||||
/// デルタのブロック種別を取得
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
match self {
|
||||
DeltaContent::Text(_) => BlockType::Text,
|
||||
DeltaContent::Thinking(_) => BlockType::Thinking,
|
||||
DeltaContent::InputJson(_) => BlockType::ToolUse,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロック停止イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockStop {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// 停止理由
|
||||
pub stop_reason: Option<StopReason>,
|
||||
}
|
||||
|
||||
impl BlockStop {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// ブロック中断イベント
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub struct BlockAbort {
|
||||
/// ブロックのインデックス
|
||||
pub index: usize,
|
||||
/// ブロックの種別
|
||||
pub block_type: BlockType,
|
||||
/// 中断理由
|
||||
pub reason: String,
|
||||
}
|
||||
|
||||
impl BlockAbort {
|
||||
pub fn block_type(&self) -> BlockType {
|
||||
self.block_type
|
||||
}
|
||||
}
|
||||
|
||||
/// 停止理由
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum StopReason {
|
||||
/// 自然終了
|
||||
EndTurn,
|
||||
/// 最大トークン数到達
|
||||
MaxTokens,
|
||||
/// ストップシーケンス到達
|
||||
StopSequence,
|
||||
/// ツール使用
|
||||
ToolUse,
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Builder / Factory helpers
|
||||
// =============================================================================
|
||||
|
||||
impl Event {
|
||||
/// テキストブロック開始イベントを作成
|
||||
pub fn text_block_start(index: usize) -> Self {
|
||||
Event::BlockStart(BlockStart {
|
||||
index,
|
||||
block_type: BlockType::Text,
|
||||
metadata: BlockMetadata::Text,
|
||||
})
|
||||
}
|
||||
|
||||
/// テキストデルタイベントを作成
|
||||
pub fn text_delta(index: usize, text: impl Into<String>) -> Self {
|
||||
Event::BlockDelta(BlockDelta {
|
||||
index,
|
||||
delta: DeltaContent::Text(text.into()),
|
||||
})
|
||||
}
|
||||
|
||||
/// テキストブロック停止イベントを作成
|
||||
pub fn text_block_stop(index: usize, stop_reason: Option<StopReason>) -> Self {
|
||||
Event::BlockStop(BlockStop {
|
||||
index,
|
||||
block_type: BlockType::Text,
|
||||
stop_reason,
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール使用ブロック開始イベントを作成
|
||||
pub fn tool_use_start(index: usize, id: impl Into<String>, name: impl Into<String>) -> Self {
|
||||
Event::BlockStart(BlockStart {
|
||||
index,
|
||||
block_type: BlockType::ToolUse,
|
||||
metadata: BlockMetadata::ToolUse {
|
||||
id: id.into(),
|
||||
name: name.into(),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール引数デルタイベントを作成
|
||||
pub fn tool_input_delta(index: usize, json: impl Into<String>) -> Self {
|
||||
Event::BlockDelta(BlockDelta {
|
||||
index,
|
||||
delta: DeltaContent::InputJson(json.into()),
|
||||
})
|
||||
}
|
||||
|
||||
/// ツール使用ブロック停止イベントを作成
|
||||
pub fn tool_use_stop(index: usize) -> Self {
|
||||
Event::BlockStop(BlockStop {
|
||||
index,
|
||||
block_type: BlockType::ToolUse,
|
||||
stop_reason: Some(StopReason::ToolUse),
|
||||
})
|
||||
}
|
||||
|
||||
/// 使用量イベントを作成
|
||||
pub fn usage(input_tokens: u64, output_tokens: u64) -> Self {
|
||||
Event::Usage(UsageEvent {
|
||||
input_tokens: Some(input_tokens),
|
||||
output_tokens: Some(output_tokens),
|
||||
total_tokens: Some(input_tokens + output_tokens),
|
||||
cache_read_input_tokens: None,
|
||||
cache_creation_input_tokens: None,
|
||||
})
|
||||
}
|
||||
|
||||
/// Pingイベントを作成
|
||||
pub fn ping() -> Self {
|
||||
Event::Ping(PingEvent { timestamp: None })
|
||||
}
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// Conversions: llm_client::event -> timeline::event
|
||||
// =============================================================================
|
||||
|
||||
impl From<crate::llm_client::event::ResponseStatus> for ResponseStatus {
|
||||
fn from(value: crate::llm_client::event::ResponseStatus) -> Self {
|
||||
match value {
|
||||
crate::llm_client::event::ResponseStatus::Started => ResponseStatus::Started,
|
||||
crate::llm_client::event::ResponseStatus::Completed => ResponseStatus::Completed,
|
||||
crate::llm_client::event::ResponseStatus::Cancelled => ResponseStatus::Cancelled,
|
||||
crate::llm_client::event::ResponseStatus::Failed => ResponseStatus::Failed,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::BlockType> for BlockType {
|
||||
fn from(value: crate::llm_client::event::BlockType) -> Self {
|
||||
match value {
|
||||
crate::llm_client::event::BlockType::Text => BlockType::Text,
|
||||
crate::llm_client::event::BlockType::Thinking => BlockType::Thinking,
|
||||
crate::llm_client::event::BlockType::ToolUse => BlockType::ToolUse,
|
||||
crate::llm_client::event::BlockType::ToolResult => BlockType::ToolResult,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::BlockMetadata> for BlockMetadata {
|
||||
fn from(value: crate::llm_client::event::BlockMetadata) -> Self {
|
||||
match value {
|
||||
crate::llm_client::event::BlockMetadata::Text => BlockMetadata::Text,
|
||||
crate::llm_client::event::BlockMetadata::Thinking => BlockMetadata::Thinking,
|
||||
crate::llm_client::event::BlockMetadata::ToolUse { id, name } => {
|
||||
BlockMetadata::ToolUse { id, name }
|
||||
}
|
||||
crate::llm_client::event::BlockMetadata::ToolResult { tool_use_id } => {
|
||||
BlockMetadata::ToolResult { tool_use_id }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::DeltaContent> for DeltaContent {
|
||||
fn from(value: crate::llm_client::event::DeltaContent) -> Self {
|
||||
match value {
|
||||
crate::llm_client::event::DeltaContent::Text(text) => DeltaContent::Text(text),
|
||||
crate::llm_client::event::DeltaContent::Thinking(text) => DeltaContent::Thinking(text),
|
||||
crate::llm_client::event::DeltaContent::InputJson(json) => {
|
||||
DeltaContent::InputJson(json)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::StopReason> for StopReason {
|
||||
fn from(value: crate::llm_client::event::StopReason) -> Self {
|
||||
match value {
|
||||
crate::llm_client::event::StopReason::EndTurn => StopReason::EndTurn,
|
||||
crate::llm_client::event::StopReason::MaxTokens => StopReason::MaxTokens,
|
||||
crate::llm_client::event::StopReason::StopSequence => StopReason::StopSequence,
|
||||
crate::llm_client::event::StopReason::ToolUse => StopReason::ToolUse,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::PingEvent> for PingEvent {
|
||||
fn from(value: crate::llm_client::event::PingEvent) -> Self {
|
||||
PingEvent {
|
||||
timestamp: value.timestamp,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::UsageEvent> for UsageEvent {
|
||||
fn from(value: crate::llm_client::event::UsageEvent) -> Self {
|
||||
UsageEvent {
|
||||
input_tokens: value.input_tokens,
|
||||
output_tokens: value.output_tokens,
|
||||
total_tokens: value.total_tokens,
|
||||
cache_read_input_tokens: value.cache_read_input_tokens,
|
||||
cache_creation_input_tokens: value.cache_creation_input_tokens,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::StatusEvent> for StatusEvent {
|
||||
fn from(value: crate::llm_client::event::StatusEvent) -> Self {
|
||||
StatusEvent {
|
||||
status: value.status.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::ErrorEvent> for ErrorEvent {
|
||||
fn from(value: crate::llm_client::event::ErrorEvent) -> Self {
|
||||
ErrorEvent {
|
||||
code: value.code,
|
||||
message: value.message,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::BlockStart> for BlockStart {
|
||||
fn from(value: crate::llm_client::event::BlockStart) -> Self {
|
||||
BlockStart {
|
||||
index: value.index,
|
||||
block_type: value.block_type.into(),
|
||||
metadata: value.metadata.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::BlockDelta> for BlockDelta {
|
||||
fn from(value: crate::llm_client::event::BlockDelta) -> Self {
|
||||
BlockDelta {
|
||||
index: value.index,
|
||||
delta: value.delta.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::BlockStop> for BlockStop {
|
||||
fn from(value: crate::llm_client::event::BlockStop) -> Self {
|
||||
BlockStop {
|
||||
index: value.index,
|
||||
block_type: value.block_type.into(),
|
||||
stop_reason: value.stop_reason.map(Into::into),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::BlockAbort> for BlockAbort {
|
||||
fn from(value: crate::llm_client::event::BlockAbort) -> Self {
|
||||
BlockAbort {
|
||||
index: value.index,
|
||||
block_type: value.block_type.into(),
|
||||
reason: value.reason,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<crate::llm_client::event::Event> for Event {
|
||||
fn from(value: crate::llm_client::event::Event) -> Self {
|
||||
match value {
|
||||
crate::llm_client::event::Event::Ping(p) => Event::Ping(p.into()),
|
||||
crate::llm_client::event::Event::Usage(u) => Event::Usage(u.into()),
|
||||
crate::llm_client::event::Event::Status(s) => Event::Status(s.into()),
|
||||
crate::llm_client::event::Event::Error(e) => Event::Error(e.into()),
|
||||
crate::llm_client::event::Event::BlockStart(s) => Event::BlockStart(s.into()),
|
||||
crate::llm_client::event::Event::BlockDelta(d) => Event::BlockDelta(d.into()),
|
||||
crate::llm_client::event::Event::BlockStop(s) => Event::BlockStop(s.into()),
|
||||
crate::llm_client::event::Event::BlockAbort(a) => Event::BlockAbort(a.into()),
|
||||
}
|
||||
}
|
||||
}
|
||||
+25
-11
@@ -9,25 +9,39 @@
|
||||
//! - [`TextBlockCollector`] - テキストブロックを収集するHandler
|
||||
//! - [`ToolCallCollector`] - ツール呼び出しを収集するHandler
|
||||
|
||||
pub mod event;
|
||||
mod text_block_collector;
|
||||
mod timeline;
|
||||
mod tool_call_collector;
|
||||
|
||||
// 公開API
|
||||
pub use event::*;
|
||||
pub use text_block_collector::TextBlockCollector;
|
||||
pub use timeline::{ErasedHandler, HandlerWrapper, Timeline};
|
||||
pub use tool_call_collector::ToolCallCollector;
|
||||
|
||||
// worker-typesからのre-export
|
||||
pub use worker_types::{
|
||||
// Core traits
|
||||
Handler, Kind,
|
||||
// Block Kinds
|
||||
TextBlockKind, ThinkingBlockKind, ToolUseBlockKind,
|
||||
// Block Events
|
||||
TextBlockEvent, TextBlockStart, TextBlockStop,
|
||||
ThinkingBlockEvent, ThinkingBlockStart, ThinkingBlockStop,
|
||||
ToolUseBlockEvent, ToolUseBlockStart, ToolUseBlockStop,
|
||||
// 型定義からのre-export
|
||||
pub use crate::handler::{
|
||||
// Meta Kinds
|
||||
ErrorKind, PingKind, StatusKind, UsageKind,
|
||||
ErrorKind,
|
||||
// Core traits
|
||||
Handler,
|
||||
Kind,
|
||||
PingKind,
|
||||
StatusKind,
|
||||
// Block Events
|
||||
TextBlockEvent,
|
||||
// Block Kinds
|
||||
TextBlockKind,
|
||||
TextBlockStart,
|
||||
TextBlockStop,
|
||||
ThinkingBlockEvent,
|
||||
ThinkingBlockKind,
|
||||
ThinkingBlockStart,
|
||||
ThinkingBlockStop,
|
||||
ToolUseBlockEvent,
|
||||
ToolUseBlockKind,
|
||||
ToolUseBlockStart,
|
||||
ToolUseBlockStop,
|
||||
UsageKind,
|
||||
};
|
||||
|
||||
@@ -3,8 +3,8 @@
|
||||
//! TimelineのTextBlockHandler として登録され、
|
||||
//! ストリーム中のテキストブロックを収集する。
|
||||
|
||||
use crate::handler::{Handler, TextBlockEvent, TextBlockKind};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use worker_types::{Handler, TextBlockEvent, TextBlockKind};
|
||||
|
||||
/// TextBlockから収集したテキスト情報を保持
|
||||
#[derive(Debug, Default)]
|
||||
@@ -85,7 +85,7 @@ impl Handler<TextBlockKind> for TextBlockCollector {
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::timeline::Timeline;
|
||||
use worker_types::Event;
|
||||
use crate::timeline::event::Event;
|
||||
|
||||
/// TextBlockCollectorが単一のテキストブロックを正しく収集することを確認
|
||||
#[test]
|
||||
|
||||
@@ -5,7 +5,8 @@
|
||||
|
||||
use std::marker::PhantomData;
|
||||
|
||||
use worker_types::*;
|
||||
use super::event::*;
|
||||
use crate::handler::*;
|
||||
|
||||
// =============================================================================
|
||||
// Type-erased Handler
|
||||
|
||||
@@ -3,8 +3,11 @@
|
||||
//! TimelineのToolUseBlockHandler として登録され、
|
||||
//! ストリーム中のToolUseブロックを収集する。
|
||||
|
||||
use crate::{
|
||||
handler::{Handler, ToolUseBlockEvent, ToolUseBlockKind},
|
||||
hook::ToolCall,
|
||||
};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use worker_types::{Handler, ToolCall, ToolUseBlockEvent, ToolUseBlockKind};
|
||||
|
||||
/// ToolUseブロックから収集したツール呼び出し情報を保持
|
||||
///
|
||||
@@ -98,7 +101,7 @@ impl Handler<ToolUseBlockKind> for ToolCallCollector {
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::timeline::Timeline;
|
||||
use worker_types::Event;
|
||||
use crate::timeline::event::Event;
|
||||
|
||||
#[test]
|
||||
fn test_collect_single_tool_call() {
|
||||
|
||||
@@ -0,0 +1,90 @@
|
||||
//! ツール定義
|
||||
//!
|
||||
//! LLMから呼び出し可能なツールを定義するためのトレイト。
|
||||
//! 通常は`#[tool]`マクロを使用して自動実装します。
|
||||
|
||||
use async_trait::async_trait;
|
||||
use serde_json::Value;
|
||||
use thiserror::Error;
|
||||
|
||||
/// ツール実行時のエラー
|
||||
#[derive(Debug, Error)]
|
||||
pub enum ToolError {
|
||||
/// 引数が不正
|
||||
#[error("Invalid argument: {0}")]
|
||||
InvalidArgument(String),
|
||||
/// 実行に失敗
|
||||
#[error("Execution failed: {0}")]
|
||||
ExecutionFailed(String),
|
||||
/// 内部エラー
|
||||
#[error("Internal error: {0}")]
|
||||
Internal(String),
|
||||
}
|
||||
|
||||
/// LLMから呼び出し可能なツールを定義するトレイト
|
||||
///
|
||||
/// ツールはLLMが外部リソースにアクセスしたり、
|
||||
/// 計算を実行したりするために使用します。
|
||||
///
|
||||
/// # 実装方法
|
||||
///
|
||||
/// 通常は`#[tool]`マクロを使用して自動実装します:
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::tool;
|
||||
///
|
||||
/// #[tool(description = "Search the web for information")]
|
||||
/// async fn search(query: String) -> String {
|
||||
/// // 検索処理
|
||||
/// format!("Results for: {}", query)
|
||||
/// }
|
||||
/// ```
|
||||
///
|
||||
/// # 手動実装
|
||||
///
|
||||
/// ```ignore
|
||||
/// use worker::tool::{Tool, ToolError};
|
||||
/// use serde_json::{json, Value};
|
||||
///
|
||||
/// struct MyTool;
|
||||
///
|
||||
/// #[async_trait::async_trait]
|
||||
/// impl Tool for MyTool {
|
||||
/// fn name(&self) -> &str { "my_tool" }
|
||||
/// fn description(&self) -> &str { "My custom tool" }
|
||||
/// fn input_schema(&self) -> Value {
|
||||
/// json!({
|
||||
/// "type": "object",
|
||||
/// "properties": {
|
||||
/// "query": { "type": "string" }
|
||||
/// },
|
||||
/// "required": ["query"]
|
||||
/// })
|
||||
/// }
|
||||
/// async fn execute(&self, input: &str) -> Result<String, ToolError> {
|
||||
/// Ok("result".to_string())
|
||||
/// }
|
||||
/// }
|
||||
/// ```
|
||||
#[async_trait]
|
||||
pub trait Tool: Send + Sync {
|
||||
/// ツール名(LLMが識別に使用)
|
||||
fn name(&self) -> &str;
|
||||
|
||||
/// ツールの説明(LLMへのプロンプトに含まれる)
|
||||
fn description(&self) -> &str;
|
||||
|
||||
/// 引数のJSON Schema
|
||||
///
|
||||
/// LLMはこのスキーマに従って引数を生成します。
|
||||
fn input_schema(&self) -> Value;
|
||||
|
||||
/// ツールを実行する
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `input_json` - LLMが生成したJSON形式の引数
|
||||
///
|
||||
/// # Returns
|
||||
/// 実行結果の文字列。この内容がLLMに返されます。
|
||||
async fn execute(&self, input_json: &str) -> Result<String, ToolError>;
|
||||
}
|
||||
+42
-43
@@ -5,15 +5,17 @@ use std::sync::{Arc, Mutex};
|
||||
use futures::StreamExt;
|
||||
use tracing::{debug, info, trace, warn};
|
||||
|
||||
use crate::timeline::{TextBlockCollector, Timeline, ToolCallCollector};
|
||||
use crate::llm_client::{ClientError, LlmClient, Request, ToolDefinition};
|
||||
use crate::subscriber_adapter::{
|
||||
ErrorSubscriberAdapter, StatusSubscriberAdapter, TextBlockSubscriberAdapter,
|
||||
ToolUseBlockSubscriberAdapter, UsageSubscriberAdapter,
|
||||
};
|
||||
use worker_types::{
|
||||
ContentPart, ControlFlow, HookError, Locked, Message, MessageContent, Mutable, Tool, ToolCall,
|
||||
ToolError, ToolResult, TurnResult, WorkerHook, WorkerState, WorkerSubscriber,
|
||||
use crate::{
|
||||
ContentPart, Message, MessageContent, Role,
|
||||
hook::{ControlFlow, HookError, ToolCall, ToolResult, TurnResult, WorkerHook},
|
||||
llm_client::{ClientError, LlmClient, Request, ToolDefinition},
|
||||
state::{Locked, Mutable, WorkerState},
|
||||
subscriber::{
|
||||
ErrorSubscriberAdapter, StatusSubscriberAdapter, TextBlockSubscriberAdapter,
|
||||
ToolUseBlockSubscriberAdapter, UsageSubscriberAdapter, WorkerSubscriber,
|
||||
},
|
||||
timeline::{TextBlockCollector, Timeline, ToolCallCollector},
|
||||
tool::{Tool, ToolError},
|
||||
};
|
||||
|
||||
// =============================================================================
|
||||
@@ -321,7 +323,7 @@ impl<C: LlmClient, S: WorkerState> Worker<C, S> {
|
||||
}
|
||||
|
||||
Some(Message {
|
||||
role: worker_types::Role::Assistant,
|
||||
role: Role::Assistant,
|
||||
content: MessageContent::Parts(parts),
|
||||
})
|
||||
}
|
||||
@@ -337,49 +339,45 @@ impl<C: LlmClient, S: WorkerState> Worker<C, S> {
|
||||
|
||||
// メッセージを追加
|
||||
for msg in &self.history {
|
||||
// worker-types::Message から llm_client::Message への変換
|
||||
// Message から llm_client::Message への変換
|
||||
request = request.message(crate::llm_client::Message {
|
||||
role: match msg.role {
|
||||
worker_types::Role::User => crate::llm_client::Role::User,
|
||||
worker_types::Role::Assistant => crate::llm_client::Role::Assistant,
|
||||
Role::User => crate::llm_client::Role::User,
|
||||
Role::Assistant => crate::llm_client::Role::Assistant,
|
||||
},
|
||||
content: match &msg.content {
|
||||
worker_types::MessageContent::Text(t) => {
|
||||
crate::llm_client::MessageContent::Text(t.clone())
|
||||
}
|
||||
worker_types::MessageContent::ToolResult {
|
||||
MessageContent::Text(t) => crate::llm_client::MessageContent::Text(t.clone()),
|
||||
MessageContent::ToolResult {
|
||||
tool_use_id,
|
||||
content,
|
||||
} => crate::llm_client::MessageContent::ToolResult {
|
||||
tool_use_id: tool_use_id.clone(),
|
||||
content: content.clone(),
|
||||
},
|
||||
worker_types::MessageContent::Parts(parts) => {
|
||||
crate::llm_client::MessageContent::Parts(
|
||||
parts
|
||||
.iter()
|
||||
.map(|p| match p {
|
||||
worker_types::ContentPart::Text { text } => {
|
||||
crate::llm_client::ContentPart::Text { text: text.clone() }
|
||||
MessageContent::Parts(parts) => crate::llm_client::MessageContent::Parts(
|
||||
parts
|
||||
.iter()
|
||||
.map(|p| match p {
|
||||
ContentPart::Text { text } => {
|
||||
crate::llm_client::ContentPart::Text { text: text.clone() }
|
||||
}
|
||||
ContentPart::ToolUse { id, name, input } => {
|
||||
crate::llm_client::ContentPart::ToolUse {
|
||||
id: id.clone(),
|
||||
name: name.clone(),
|
||||
input: input.clone(),
|
||||
}
|
||||
worker_types::ContentPart::ToolUse { id, name, input } => {
|
||||
crate::llm_client::ContentPart::ToolUse {
|
||||
id: id.clone(),
|
||||
name: name.clone(),
|
||||
input: input.clone(),
|
||||
}
|
||||
}
|
||||
worker_types::ContentPart::ToolResult {
|
||||
tool_use_id,
|
||||
content,
|
||||
} => crate::llm_client::ContentPart::ToolResult {
|
||||
tool_use_id: tool_use_id.clone(),
|
||||
content: content.clone(),
|
||||
},
|
||||
})
|
||||
.collect(),
|
||||
)
|
||||
}
|
||||
}
|
||||
ContentPart::ToolResult {
|
||||
tool_use_id,
|
||||
content,
|
||||
} => crate::llm_client::ContentPart::ToolResult {
|
||||
tool_use_id: tool_use_id.clone(),
|
||||
content: content.clone(),
|
||||
},
|
||||
})
|
||||
.collect(),
|
||||
),
|
||||
},
|
||||
});
|
||||
}
|
||||
@@ -550,7 +548,8 @@ impl<C: LlmClient, S: WorkerState> Worker<C, S> {
|
||||
}
|
||||
}
|
||||
let event = event_result?;
|
||||
self.timeline.dispatch(&event);
|
||||
let timeline_event: crate::timeline::event::Event = event.into();
|
||||
self.timeline.dispatch(&timeline_event);
|
||||
}
|
||||
debug!(event_count = event_count, "Stream completed");
|
||||
|
||||
|
||||
@@ -8,9 +8,9 @@ use std::sync::{Arc, Mutex};
|
||||
|
||||
use async_trait::async_trait;
|
||||
use futures::Stream;
|
||||
use worker::llm_client::event::{BlockType, DeltaContent, Event};
|
||||
use worker::llm_client::{ClientError, LlmClient, Request};
|
||||
use worker::timeline::{Handler, TextBlockEvent, TextBlockKind, Timeline};
|
||||
use worker_types::{BlockType, DeltaContent, Event};
|
||||
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
@@ -267,7 +267,8 @@ pub fn assert_timeline_integration(subdir: &str) {
|
||||
});
|
||||
|
||||
for event in &events {
|
||||
timeline.dispatch(event);
|
||||
let timeline_event: worker::timeline::event::Event = event.clone().into();
|
||||
timeline.dispatch(&timeline_event);
|
||||
}
|
||||
|
||||
let texts = collected.lock().unwrap();
|
||||
|
||||
@@ -8,10 +8,9 @@ use std::time::{Duration, Instant};
|
||||
|
||||
use async_trait::async_trait;
|
||||
use worker::Worker;
|
||||
use worker_types::{
|
||||
ControlFlow, Event, HookError, ResponseStatus, StatusEvent, Tool, ToolCall, ToolError,
|
||||
ToolResult, WorkerHook,
|
||||
};
|
||||
use worker::hook::{ControlFlow, HookError, ToolCall, ToolResult, WorkerHook};
|
||||
use worker::llm_client::event::{Event, ResponseStatus, StatusEvent};
|
||||
use worker::tool::{Tool, ToolError};
|
||||
|
||||
mod common;
|
||||
use common::MockLlmClient;
|
||||
|
||||
@@ -7,12 +7,12 @@ mod common;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use common::MockLlmClient;
|
||||
use worker::subscriber::WorkerSubscriber;
|
||||
use worker::Worker;
|
||||
use worker_types::{
|
||||
ErrorEvent, Event, ResponseStatus, StatusEvent, TextBlockEvent, ToolCall, ToolUseBlockEvent,
|
||||
UsageEvent,
|
||||
};
|
||||
use worker::hook::ToolCall;
|
||||
use worker::llm_client::event::{Event, ResponseStatus, StatusEvent as ClientStatusEvent};
|
||||
use worker::subscriber::WorkerSubscriber;
|
||||
use worker::timeline::event::{ErrorEvent, StatusEvent, UsageEvent};
|
||||
use worker::timeline::{TextBlockEvent, ToolUseBlockEvent};
|
||||
|
||||
// =============================================================================
|
||||
// Test Subscriber
|
||||
@@ -101,7 +101,7 @@ async fn test_subscriber_text_block_events() {
|
||||
Event::text_delta(0, "Hello, "),
|
||||
Event::text_delta(0, "World!"),
|
||||
Event::text_block_stop(0, None),
|
||||
Event::Status(StatusEvent {
|
||||
Event::Status(ClientStatusEvent {
|
||||
status: ResponseStatus::Completed,
|
||||
}),
|
||||
];
|
||||
@@ -141,7 +141,7 @@ async fn test_subscriber_tool_call_complete() {
|
||||
Event::tool_input_delta(0, r#"{"city":"#),
|
||||
Event::tool_input_delta(0, r#""Tokyo"}"#),
|
||||
Event::tool_use_stop(0),
|
||||
Event::Status(StatusEvent {
|
||||
Event::Status(ClientStatusEvent {
|
||||
status: ResponseStatus::Completed,
|
||||
}),
|
||||
];
|
||||
@@ -172,7 +172,7 @@ async fn test_subscriber_turn_events() {
|
||||
Event::text_block_start(0),
|
||||
Event::text_delta(0, "Done!"),
|
||||
Event::text_block_stop(0, None),
|
||||
Event::Status(StatusEvent {
|
||||
Event::Status(ClientStatusEvent {
|
||||
status: ResponseStatus::Completed,
|
||||
}),
|
||||
];
|
||||
@@ -210,7 +210,7 @@ async fn test_subscriber_usage_events() {
|
||||
Event::text_delta(0, "Hello"),
|
||||
Event::text_block_stop(0, None),
|
||||
Event::usage(100, 50),
|
||||
Event::Status(StatusEvent {
|
||||
Event::Status(ClientStatusEvent {
|
||||
status: ResponseStatus::Completed,
|
||||
}),
|
||||
];
|
||||
|
||||
@@ -9,8 +9,8 @@ use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use schemars;
|
||||
use serde;
|
||||
|
||||
use worker::tool::Tool;
|
||||
use worker_macros::tool_registry;
|
||||
use worker_types::Tool;
|
||||
|
||||
// =============================================================================
|
||||
// Test: Basic Tool Generation
|
||||
|
||||
@@ -12,7 +12,7 @@ use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use async_trait::async_trait;
|
||||
use common::MockLlmClient;
|
||||
use worker::Worker;
|
||||
use worker_types::{Tool, ToolError};
|
||||
use worker::tool::{Tool, ToolError};
|
||||
|
||||
/// フィクスチャディレクトリのパス
|
||||
fn fixtures_dir() -> std::path::PathBuf {
|
||||
@@ -100,7 +100,7 @@ fn test_mock_client_from_fixture() {
|
||||
/// fixtureファイルを使わず、プログラムでイベントを構築してクライアントを作成する。
|
||||
#[test]
|
||||
fn test_mock_client_from_events() {
|
||||
use worker_types::Event;
|
||||
use worker::llm_client::event::Event;
|
||||
|
||||
// 直接イベントを指定
|
||||
let events = vec![
|
||||
@@ -178,7 +178,7 @@ async fn test_worker_tool_call() {
|
||||
/// テストの独立性を高め、外部ファイルへの依存を排除したい場合に有用。
|
||||
#[tokio::test]
|
||||
async fn test_worker_with_programmatic_events() {
|
||||
use worker_types::{Event, ResponseStatus, StatusEvent};
|
||||
use worker::llm_client::event::{Event, ResponseStatus, StatusEvent};
|
||||
|
||||
// プログラムでイベントシーケンスを構築
|
||||
let events = vec![
|
||||
@@ -205,8 +205,8 @@ async fn test_worker_with_programmatic_events() {
|
||||
/// id, name, input(JSON)を正しく抽出できることを検証する。
|
||||
#[tokio::test]
|
||||
async fn test_tool_call_collector_integration() {
|
||||
use worker::llm_client::event::Event;
|
||||
use worker::timeline::{Timeline, ToolCallCollector};
|
||||
use worker_types::Event;
|
||||
|
||||
// ToolUseブロックを含むイベントシーケンス
|
||||
let events = vec![
|
||||
@@ -222,7 +222,8 @@ async fn test_tool_call_collector_integration() {
|
||||
|
||||
// イベントをディスパッチ
|
||||
for event in &events {
|
||||
timeline.dispatch(event);
|
||||
let timeline_event: worker::timeline::event::Event = event.clone().into();
|
||||
timeline.dispatch(&timeline_event);
|
||||
}
|
||||
|
||||
// 収集されたToolCallを確認
|
||||
|
||||
@@ -7,7 +7,8 @@ mod common;
|
||||
|
||||
use common::MockLlmClient;
|
||||
use worker::Worker;
|
||||
use worker_types::{Event, Message, MessageContent, ResponseStatus, StatusEvent};
|
||||
use worker::llm_client::event::{Event, ResponseStatus, StatusEvent};
|
||||
use worker::{Message, MessageContent};
|
||||
|
||||
// =============================================================================
|
||||
// Mutable状態のテスト
|
||||
|
||||
Reference in New Issue
Block a user