use harmony_response::{
chat::{Conversation, Message, Role, SystemContent},
load_harmony_encoding, HarmonyEncodingName, StreamableParser,
};
use std::{thread, time::Duration};
fn main() -> anyhow::Result<()> {
println!("🌊 Harmony Response Format - Streaming Parser Example");
println!("=====================================================");
println!("\n🔧 Step 1: Loading encoding and creating sample data");
match load_harmony_encoding(HarmonyEncodingName::HarmonyGptOss) {
Ok(encoding) => {
println!("✅ Successfully loaded encoding: {}", encoding.name());
let conversation = Conversation::from_messages([
Message::from_role_and_content(
Role::System,
SystemContent::new()
.with_model_identity("Streaming Demo Assistant")
.with_required_channels(["analysis", "commentary", "final"])
),
Message::from_role_and_content(Role::User, "Explain quantum computing in simple terms."),
]);
let completion_tokens = encoding
.render_conversation_for_completion(&conversation, Role::Assistant, None)?;
println!("✓ Created sample with {} tokens to simulate streaming from OpenAI", completion_tokens.len());
println!(" (In real usage, these tokens would come from OpenAI's streaming API)");
println!("\n🎯 Step 2: Creating streaming parser");
let mut parser = StreamableParser::new(encoding.clone(), Some(Role::Assistant))?;
println!("✓ Created parser expecting Assistant role");
println!("\n📡 Step 3: Simulating token streaming");
println!("Processing tokens in real-time...\n");
for (i, &token) in completion_tokens.iter().enumerate() {
parser.process(token)?;
let current_role = parser.current_role();
let current_channel = parser.current_channel();
let current_recipient = parser.current_recipient();
let delta = parser.last_content_delta()?;
print!("[{}] Token {}: {} ", i + 1, token, token);
if let Some(role) = current_role {
print!("(Role: {}) ", role);
}
if let Some(channel) = ¤t_channel {
print!("(Channel: {}) ", channel);
}
if let Some(recipient) = ¤t_recipient {
print!("(To: {}) ", recipient);
}
if let Some(delta_content) = &delta {
if !delta_content.is_empty() {
print!("→ \"{}\"", delta_content.replace('\n', "\\n"));
}
} else {
print!("(incomplete UTF-8)");
}
println!();
if i % 10 == 9 {
if let Ok(current_content) = parser.current_content() {
if !current_content.is_empty() {
println!(" Current content: \"{}\"",
current_content.replace('\n', "\\n").chars().take(50).collect::<String>()
+ if current_content.len() > 50 { "..." } else { "" });
}
}
println!();
}
thread::sleep(Duration::from_millis(50));
}
println!("\n🏁 Step 4: Finalizing parsing");
parser.process_eos()?;
println!("✓ Processed end-of-stream");
println!("\n📋 Step 5: Extracting parsed messages");
let messages = parser.into_messages();
println!("✓ Extracted {} messages from stream", messages.len());
for (i, message) in messages.iter().enumerate() {
println!("\n Message {}:", i + 1);
println!(" Role: {}", message.author.role);
if let Some(channel) = &message.channel {
println!(" Channel: {}", channel);
}
if let Some(recipient) = &message.recipient {
println!(" Recipient: {}", recipient);
}
for (j, content) in message.content.iter().enumerate() {
match content {
harmony_response::chat::Content::Text(text) => {
let preview = text.text.chars().take(100).collect::<String>();
println!(" Content {}: \"{}{}\"",
j + 1,
preview.replace('\n', "\\n"),
if text.text.len() > 100 { "..." } else { "" }
);
}
_ => println!(" Content {}: [Non-text content]", j + 1),
}
}
}
println!("\n🎉 Streaming parser example completed successfully!");
}
Err(e) => {
println!("⚠️ Encoding load failed (expected without network access)");
println!(" Error: {}", e);
println!("\n🔄 Running demo with mock tokens...");
let mock_tokens = vec![
200006, 15847, 200008, 9906, 11, 9906, 0, 200007, ];
println!("Using mock token sequence: {:?}", mock_tokens);
println!("Note: These tokens are for demonstration only.");
println!("In a real scenario, the encoding would be loaded from the network.");
}
}
Ok(())
}