|
@@ -0,0 +1,862 @@
|
|
|
|
|
+#include "smartbotic/llm/event/event_processor.hpp"
|
|
|
|
|
+#include "smartbotic/llm/provider/provider_factory.hpp"
|
|
|
|
|
+#include "smartbotic/llm/session/session_store.hpp"
|
|
|
|
|
+#include "smartbotic/llm/tool/tool_registry.hpp"
|
|
|
|
|
+
|
|
|
|
|
+#include <chrono>
|
|
|
|
|
+#include <spdlog/spdlog.h>
|
|
|
|
|
+#include <uuid.h>
|
|
|
|
|
+
|
|
|
|
|
+namespace smartbotic::llm::event {
|
|
|
|
|
+
|
|
|
|
|
+namespace {
|
|
|
|
|
+
|
|
|
|
|
+// Generate a unique ID
|
|
|
|
|
+auto GenerateId() -> std::string {
|
|
|
|
|
+ std::random_device rd;
|
|
|
|
|
+ auto seed_data = std::array<int, std::mt19937::state_size>{};
|
|
|
|
|
+ std::generate(std::begin(seed_data), std::end(seed_data), std::ref(rd));
|
|
|
|
|
+ std::seed_seq seq(std::begin(seed_data), std::end(seed_data));
|
|
|
|
|
+ std::mt19937 generator(seq);
|
|
|
|
|
+ uuids::uuid_random_generator gen{generator};
|
|
|
|
|
+ return uuids::to_string(gen());
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// Get current timestamp as ISO8601 string
|
|
|
|
|
+auto GetCurrentTimestamp() -> std::string {
|
|
|
|
|
+ auto now = std::chrono::system_clock::now();
|
|
|
|
|
+ auto time_t_now = std::chrono::system_clock::to_time_t(now);
|
|
|
|
|
+ std::stringstream ss;
|
|
|
|
|
+ ss << std::put_time(std::gmtime(&time_t_now), "%FT%TZ");
|
|
|
|
|
+ return ss.str();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// Convert timestamp string to proto Timestamp
|
|
|
|
|
+auto StringToTimestamp(const std::string& str) -> ::smartbotic::llm::Timestamp {
|
|
|
|
|
+ ::smartbotic::llm::Timestamp ts;
|
|
|
|
|
+ // Parse ISO8601 format - simplified for now
|
|
|
|
|
+ std::tm tm = {};
|
|
|
|
|
+ std::istringstream ss(str);
|
|
|
|
|
+ ss >> std::get_time(&tm, "%Y-%m-%dT%H:%M:%S");
|
|
|
|
|
+ if (!ss.fail()) {
|
|
|
|
|
+ auto time = std::mktime(&tm);
|
|
|
|
|
+ ts.set_seconds(time);
|
|
|
|
|
+ }
|
|
|
|
|
+ return ts;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+// Convert proto Timestamp to string
|
|
|
|
|
+auto TimestampToString(const ::smartbotic::llm::Timestamp& ts) -> std::string {
|
|
|
|
|
+ auto time_t_val = static_cast<time_t>(ts.seconds());
|
|
|
|
|
+ std::stringstream ss;
|
|
|
|
|
+ ss << std::put_time(std::gmtime(&time_t_val), "%FT%TZ");
|
|
|
|
|
+ return ss.str();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+} // namespace
|
|
|
|
|
+
|
|
|
|
|
+EventProcessor::EventProcessor(
|
|
|
|
|
+ const EventProcessorConfig& config,
|
|
|
|
|
+ std::shared_ptr<provider::ProviderFactory> provider_factory,
|
|
|
|
|
+ std::shared_ptr<session::ISessionStore> session_store,
|
|
|
|
|
+ std::shared_ptr<tool::ToolRegistry> tool_registry
|
|
|
|
|
+) : config_(config),
|
|
|
|
|
+ provider_factory_(std::move(provider_factory)),
|
|
|
|
|
+ session_store_(std::move(session_store)),
|
|
|
|
|
+ tool_registry_(std::move(tool_registry)) {
|
|
|
|
|
+
|
|
|
|
|
+ // Generate node ID if not provided
|
|
|
|
|
+ if (config_.node_id.empty()) {
|
|
|
|
|
+ config_.node_id = GenerateId();
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+EventProcessor::~EventProcessor() {
|
|
|
|
|
+ Stop();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::Connect() -> bool {
|
|
|
|
|
+ try {
|
|
|
|
|
+ channel_ = grpc::CreateChannel(
|
|
|
|
|
+ config_.database_address,
|
|
|
|
|
+ grpc::InsecureChannelCredentials()
|
|
|
|
|
+ );
|
|
|
|
|
+
|
|
|
|
|
+ document_stub_ = ::smartbotic::database::DocumentService::NewStub(channel_);
|
|
|
|
|
+ collection_stub_ = ::smartbotic::database::CollectionService::NewStub(channel_);
|
|
|
|
|
+ query_stub_ = ::smartbotic::database::QueryService::NewStub(channel_);
|
|
|
|
|
+ subscription_stub_ = ::smartbotic::database::SubscriptionService::NewStub(channel_);
|
|
|
|
|
+
|
|
|
|
|
+ // Wait for connection
|
|
|
|
|
+ auto deadline = std::chrono::system_clock::now() + std::chrono::seconds(10);
|
|
|
|
|
+ if (!channel_->WaitForConnected(deadline)) {
|
|
|
|
|
+ spdlog::error("EventProcessor: Failed to connect to database at {}",
|
|
|
|
|
+ config_.database_address);
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("EventProcessor: Connected to database at {}", config_.database_address);
|
|
|
|
|
+ return true;
|
|
|
|
|
+ } catch (const std::exception& e) {
|
|
|
|
|
+ spdlog::error("EventProcessor: Exception connecting to database: {}", e.what());
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::Initialize() -> bool {
|
|
|
|
|
+ if (initialized_) {
|
|
|
|
|
+ return true;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!Connect()) {
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Create collections if they don't exist
|
|
|
|
|
+ grpc::ClientContext events_ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateCollectionRequest events_req;
|
|
|
|
|
+ events_req.set_name(config_.events_collection);
|
|
|
|
|
+ ::smartbotic::database::CreateCollectionResponse events_resp;
|
|
|
|
|
+
|
|
|
|
|
+ auto events_status = collection_stub_->CreateCollection(&events_ctx, events_req, &events_resp);
|
|
|
|
|
+ if (!events_status.ok() &&
|
|
|
|
|
+ events_status.error_code() != grpc::StatusCode::ALREADY_EXISTS) {
|
|
|
|
|
+ spdlog::error("EventProcessor: Failed to create events collection: {}",
|
|
|
|
|
+ events_status.error_message());
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ spdlog::info("EventProcessor: Events collection {} ready", config_.events_collection);
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext chunks_ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateCollectionRequest chunks_req;
|
|
|
|
|
+ chunks_req.set_name(config_.chunks_collection);
|
|
|
|
|
+ ::smartbotic::database::CreateCollectionResponse chunks_resp;
|
|
|
|
|
+
|
|
|
|
|
+ auto chunks_status = collection_stub_->CreateCollection(&chunks_ctx, chunks_req, &chunks_resp);
|
|
|
|
|
+ if (!chunks_status.ok() &&
|
|
|
|
|
+ chunks_status.error_code() != grpc::StatusCode::ALREADY_EXISTS) {
|
|
|
|
|
+ spdlog::error("EventProcessor: Failed to create chunks collection: {}",
|
|
|
|
|
+ chunks_status.error_message());
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ spdlog::info("EventProcessor: Chunks collection {} ready", config_.chunks_collection);
|
|
|
|
|
+
|
|
|
|
|
+ initialized_ = true;
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::Start() {
|
|
|
|
|
+ if (running_) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!initialized_ && !Initialize()) {
|
|
|
|
|
+ spdlog::error("EventProcessor: Cannot start - initialization failed");
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ running_ = true;
|
|
|
|
|
+ processing_thread_ = std::thread([this]() {
|
|
|
|
|
+ ProcessingLoop();
|
|
|
|
|
+ });
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("EventProcessor: Started processing loop (node_id={})", config_.node_id);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::Stop() {
|
|
|
|
|
+ if (!running_) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ running_ = false;
|
|
|
|
|
+
|
|
|
|
|
+ if (processing_thread_.joinable()) {
|
|
|
|
|
+ processing_thread_.join();
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::info("EventProcessor: Stopped processing loop");
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::IsRunning() const -> bool {
|
|
|
|
|
+ return running_;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::ProcessingLoop() {
|
|
|
|
|
+ while (running_) {
|
|
|
|
|
+ try {
|
|
|
|
|
+ // Recover stale events periodically
|
|
|
|
|
+ RecoverStaleEvents();
|
|
|
|
|
+
|
|
|
|
|
+ // Poll for pending events
|
|
|
|
|
+ auto events = PollPendingEvents();
|
|
|
|
|
+
|
|
|
|
|
+ for (const auto& event : events) {
|
|
|
|
|
+ if (!running_) break;
|
|
|
|
|
+
|
|
|
|
|
+ // Try to claim the event
|
|
|
|
|
+ if (ClaimEvent(event.id())) {
|
|
|
|
|
+ // Process the event
|
|
|
|
|
+ ProcessEvent(event);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Sleep between polls
|
|
|
|
|
+ std::this_thread::sleep_for(
|
|
|
|
|
+ std::chrono::milliseconds(config_.poll_interval_ms)
|
|
|
|
|
+ );
|
|
|
|
|
+ } catch (const std::exception& e) {
|
|
|
|
|
+ spdlog::error("EventProcessor: Error in processing loop: {}", e.what());
|
|
|
|
|
+ std::this_thread::sleep_for(std::chrono::seconds(1));
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::PollPendingEvents() -> std::vector<::smartbotic::llm::MessageEvent> {
|
|
|
|
|
+ std::vector<::smartbotic::llm::MessageEvent> events;
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::QueryRequest req;
|
|
|
|
|
+ req.set_collection(config_.events_collection);
|
|
|
|
|
+
|
|
|
|
|
+ // Query for pending events
|
|
|
|
|
+ auto* query = req.mutable_query();
|
|
|
|
|
+ auto* status_filter = query->add_filters();
|
|
|
|
|
+ status_filter->set_field("status");
|
|
|
|
|
+ status_filter->set_operator_(::smartbotic::database::FilterOperator::FILTER_OPERATOR_EQUALS);
|
|
|
|
|
+ auto* status_val = status_filter->mutable_value();
|
|
|
|
|
+ status_val->set_int64_value(
|
|
|
|
|
+ static_cast<int64_t>(::smartbotic::llm::MESSAGE_EVENT_STATUS_PENDING)
|
|
|
|
|
+ );
|
|
|
|
|
+
|
|
|
|
|
+ req.set_limit(config_.max_concurrent_events);
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::QueryResponse resp;
|
|
|
|
|
+ auto status = query_stub_->Query(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ spdlog::warn("EventProcessor: Failed to query pending events: {}",
|
|
|
|
|
+ status.error_message());
|
|
|
|
|
+ return events;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ for (const auto& doc : resp.documents()) {
|
|
|
|
|
+ events.push_back(DocumentToEvent(doc.data(), doc.id()));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return events;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::ClaimEvent(const std::string& event_id) -> bool {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(processing_mutex_);
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentRequest req;
|
|
|
|
|
+ req.set_collection(config_.events_collection);
|
|
|
|
|
+ req.set_id(event_id);
|
|
|
|
|
+
|
|
|
|
|
+ // Only claim if still pending (optimistic locking via version check would be better)
|
|
|
|
|
+ auto* data = req.mutable_data();
|
|
|
|
|
+ auto* fields = data->mutable_fields();
|
|
|
|
|
+
|
|
|
|
|
+ // Set status to processing
|
|
|
|
|
+ (*fields)["status"].set_int64_value(
|
|
|
|
|
+ static_cast<int64_t>(::smartbotic::llm::MESSAGE_EVENT_STATUS_PROCESSING)
|
|
|
|
|
+ );
|
|
|
|
|
+ (*fields)["processing_node_id"].set_string_value(config_.node_id);
|
|
|
|
|
+ (*fields)["started_at"].set_string_value(GetCurrentTimestamp());
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentResponse resp;
|
|
|
|
|
+ auto status = document_stub_->UpdateDocument(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ // Another node may have claimed it
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ spdlog::debug("EventProcessor: Claimed event {}", event_id);
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::ProcessEvent(const ::smartbotic::llm::MessageEvent& event) {
|
|
|
|
|
+ spdlog::info("EventProcessor: Processing event {} (type={})",
|
|
|
|
|
+ event.id(), static_cast<int>(event.event_type()));
|
|
|
|
|
+
|
|
|
|
|
+ try {
|
|
|
|
|
+ // Get the session
|
|
|
|
|
+ auto session_result = session_store_->GetSession(
|
|
|
|
|
+ event.session_id(), event.user_id(), true
|
|
|
|
|
+ );
|
|
|
|
|
+ if (!session_result.ok) {
|
|
|
|
|
+ UpdateEventStatus(event.id(), ::smartbotic::llm::MESSAGE_EVENT_STATUS_FAILED,
|
|
|
|
|
+ "Failed to get session: " + session_result.error);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto& session = session_result.value;
|
|
|
|
|
+
|
|
|
|
|
+ // Update session streaming state
|
|
|
|
|
+ ::smartbotic::llm::StreamingState streaming_state;
|
|
|
|
|
+ streaming_state.set_is_streaming(true);
|
|
|
|
|
+ streaming_state.set_current_event_id(event.id());
|
|
|
|
|
+ streaming_state.set_last_chunk_sequence(0);
|
|
|
|
|
+ *streaming_state.mutable_started_at() = StringToTimestamp(GetCurrentTimestamp());
|
|
|
|
|
+ UpdateSessionStreamingState(event.session_id(), streaming_state);
|
|
|
|
|
+
|
|
|
|
|
+ // Build the chat request based on event type
|
|
|
|
|
+ // Note: Actual provider call logic would go here
|
|
|
|
|
+ // For now, we'll simulate the streaming process
|
|
|
|
|
+
|
|
|
|
|
+ std::string message_id = GenerateId();
|
|
|
|
|
+ int32_t sequence = 0;
|
|
|
|
|
+
|
|
|
|
|
+ // Get provider
|
|
|
|
|
+ std::string provider_id = event.provider_id();
|
|
|
|
|
+ if (provider_id.empty()) {
|
|
|
|
|
+ provider_id = session.model_id(); // Use session's model to determine provider
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto provider = provider_factory_->GetProvider(provider_id);
|
|
|
|
|
+ if (!provider) {
|
|
|
|
|
+ UpdateEventStatus(event.id(), ::smartbotic::llm::MESSAGE_EVENT_STATUS_FAILED,
|
|
|
|
|
+ "Provider not found: " + provider_id);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Build ChatRequest from session and event
|
|
|
|
|
+ provider::ChatRequest chat_request;
|
|
|
|
|
+ chat_request.model = event.model_id().empty() ? session.model_id() : event.model_id();
|
|
|
|
|
+
|
|
|
|
|
+ // Convert session messages to provider format
|
|
|
|
|
+ for (const auto& msg : session.messages()) {
|
|
|
|
|
+ provider::Message prov_msg;
|
|
|
|
|
+ prov_msg.role = static_cast<provider::MessageRole>(msg.role());
|
|
|
|
|
+ for (const auto& part : msg.content()) {
|
|
|
|
|
+ if (part.has_text()) {
|
|
|
|
|
+ prov_msg.content.push_back(provider::TextContent{part.text()});
|
|
|
|
|
+ }
|
|
|
|
|
+ // Handle other content types as needed
|
|
|
|
|
+ }
|
|
|
|
|
+ chat_request.messages.push_back(std::move(prov_msg));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Add user message or tool results based on event type
|
|
|
|
|
+ if (event.event_type() == ::smartbotic::llm::MESSAGE_EVENT_TYPE_USER_MESSAGE &&
|
|
|
|
|
+ event.has_user_message()) {
|
|
|
|
|
+ provider::Message user_msg;
|
|
|
|
|
+ user_msg.role = provider::MessageRole::kUser;
|
|
|
|
|
+ for (const auto& part : event.user_message().content()) {
|
|
|
|
|
+ if (part.has_text()) {
|
|
|
|
|
+ user_msg.content.push_back(provider::TextContent{part.text()});
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ chat_request.messages.push_back(std::move(user_msg));
|
|
|
|
|
+ } else if (event.event_type() == ::smartbotic::llm::MESSAGE_EVENT_TYPE_TOOL_RESULT) {
|
|
|
|
|
+ for (const auto& result : event.tool_results()) {
|
|
|
|
|
+ provider::Message tool_msg;
|
|
|
|
|
+ tool_msg.role = provider::MessageRole::kTool;
|
|
|
|
|
+ tool_msg.content.push_back(provider::ToolResultContent{
|
|
|
|
|
+ result.tool_call_id(),
|
|
|
|
|
+ result.content(),
|
|
|
|
|
+ result.is_error()
|
|
|
|
|
+ });
|
|
|
|
|
+ chat_request.messages.push_back(std::move(tool_msg));
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Convert tools from event
|
|
|
|
|
+ for (const auto& tool_def : event.tools()) {
|
|
|
|
|
+ provider::ToolDefinition tool;
|
|
|
|
|
+ tool.name = tool_def.name();
|
|
|
|
|
+ tool.description = tool_def.description();
|
|
|
|
|
+ tool.input_schema = tool_def.input_schema();
|
|
|
|
|
+ chat_request.tools.push_back(std::move(tool));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Stream callback to write chunks
|
|
|
|
|
+ auto stream_callback = [&](const provider::StreamChunk& chunk) -> bool {
|
|
|
|
|
+ ::smartbotic::llm::ResponseChunk response_chunk;
|
|
|
|
|
+ response_chunk.set_id(GenerateId());
|
|
|
|
|
+ response_chunk.set_workspace_id(event.workspace_id());
|
|
|
|
|
+ response_chunk.set_session_id(event.session_id());
|
|
|
|
|
+ response_chunk.set_message_event_id(event.id());
|
|
|
|
|
+ response_chunk.set_message_id(message_id);
|
|
|
|
|
+ response_chunk.set_sequence(sequence++);
|
|
|
|
|
+ *response_chunk.mutable_created_at() = StringToTimestamp(GetCurrentTimestamp());
|
|
|
|
|
+
|
|
|
|
|
+ if (!chunk.content_delta.empty()) {
|
|
|
|
|
+ response_chunk.set_chunk_type(::smartbotic::llm::RESPONSE_CHUNK_TYPE_CONTENT);
|
|
|
|
|
+ response_chunk.set_content_delta(chunk.content_delta);
|
|
|
|
|
+ } else if (!chunk.thinking_delta.empty()) {
|
|
|
|
|
+ response_chunk.set_chunk_type(::smartbotic::llm::RESPONSE_CHUNK_TYPE_THINKING);
|
|
|
|
|
+ response_chunk.set_thinking_delta(chunk.thinking_delta);
|
|
|
|
|
+ } else if (chunk.tool_call.has_value()) {
|
|
|
|
|
+ response_chunk.set_chunk_type(::smartbotic::llm::RESPONSE_CHUNK_TYPE_TOOL_CALL);
|
|
|
|
|
+ auto* tc = response_chunk.mutable_tool_call();
|
|
|
|
|
+ tc->set_id(chunk.tool_call->id);
|
|
|
|
|
+ tc->set_name(chunk.tool_call->name);
|
|
|
|
|
+ tc->set_arguments(chunk.tool_call->arguments);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ WriteChunk(response_chunk);
|
|
|
|
|
+
|
|
|
|
|
+ // Update streaming state
|
|
|
|
|
+ streaming_state.set_last_chunk_sequence(sequence);
|
|
|
|
|
+ UpdateSessionStreamingState(event.session_id(), streaming_state);
|
|
|
|
|
+
|
|
|
|
|
+ return true; // Continue streaming
|
|
|
|
|
+ };
|
|
|
|
|
+
|
|
|
|
|
+ // Call the provider
|
|
|
|
|
+ auto result = provider->ChatStream(chat_request, stream_callback);
|
|
|
|
|
+
|
|
|
|
|
+ // Write final metadata chunk
|
|
|
|
|
+ ::smartbotic::llm::ResponseChunk final_chunk;
|
|
|
|
|
+ final_chunk.set_id(GenerateId());
|
|
|
|
|
+ final_chunk.set_workspace_id(event.workspace_id());
|
|
|
|
|
+ final_chunk.set_session_id(event.session_id());
|
|
|
|
|
+ final_chunk.set_message_event_id(event.id());
|
|
|
|
|
+ final_chunk.set_message_id(message_id);
|
|
|
|
|
+ final_chunk.set_sequence(sequence);
|
|
|
|
|
+ final_chunk.set_chunk_type(::smartbotic::llm::RESPONSE_CHUNK_TYPE_METADATA);
|
|
|
|
|
+ final_chunk.set_is_final(true);
|
|
|
|
|
+ *final_chunk.mutable_created_at() = StringToTimestamp(GetCurrentTimestamp());
|
|
|
|
|
+
|
|
|
|
|
+ if (result.ok) {
|
|
|
|
|
+ auto* metadata = final_chunk.mutable_metadata();
|
|
|
|
|
+ auto* usage = metadata->mutable_usage();
|
|
|
|
|
+ usage->set_prompt_tokens(result.value.usage.prompt_tokens);
|
|
|
|
|
+ usage->set_completion_tokens(result.value.usage.completion_tokens);
|
|
|
|
|
+ usage->set_total_tokens(result.value.usage.total_tokens);
|
|
|
|
|
+ metadata->set_finish_reason(
|
|
|
|
|
+ static_cast<::smartbotic::llm::FinishReason>(result.value.finish_reason)
|
|
|
|
|
+ );
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ WriteChunk(final_chunk);
|
|
|
|
|
+
|
|
|
|
|
+ // Clear streaming state
|
|
|
|
|
+ streaming_state.set_is_streaming(false);
|
|
|
|
|
+ streaming_state.set_last_chunk_sequence(sequence);
|
|
|
|
|
+ UpdateSessionStreamingState(event.session_id(), streaming_state);
|
|
|
|
|
+
|
|
|
|
|
+ // Update event status
|
|
|
|
|
+ if (result.ok) {
|
|
|
|
|
+ UpdateEventStatus(event.id(), ::smartbotic::llm::MESSAGE_EVENT_STATUS_COMPLETED);
|
|
|
|
|
+ } else {
|
|
|
|
|
+ UpdateEventStatus(event.id(), ::smartbotic::llm::MESSAGE_EVENT_STATUS_FAILED,
|
|
|
|
|
+ result.error);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ } catch (const std::exception& e) {
|
|
|
|
|
+ spdlog::error("EventProcessor: Exception processing event {}: {}",
|
|
|
|
|
+ event.id(), e.what());
|
|
|
|
|
+ UpdateEventStatus(event.id(), ::smartbotic::llm::MESSAGE_EVENT_STATUS_FAILED,
|
|
|
|
|
+ e.what());
|
|
|
|
|
+
|
|
|
|
|
+ // Clear streaming state on error
|
|
|
|
|
+ ::smartbotic::llm::StreamingState streaming_state;
|
|
|
|
|
+ streaming_state.set_is_streaming(false);
|
|
|
|
|
+ UpdateSessionStreamingState(event.session_id(), streaming_state);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::WriteChunk(const ::smartbotic::llm::ResponseChunk& chunk) {
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateDocumentRequest req;
|
|
|
|
|
+ req.set_collection(config_.chunks_collection);
|
|
|
|
|
+ *req.mutable_data() = ChunkToDocument(chunk);
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::CreateDocumentResponse resp;
|
|
|
|
|
+ auto status = document_stub_->CreateDocument(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ spdlog::warn("EventProcessor: Failed to write chunk: {}", status.error_message());
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::UpdateEventStatus(const std::string& event_id,
|
|
|
|
|
+ ::smartbotic::llm::MessageEventStatus status,
|
|
|
|
|
+ const std::string& error) {
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentRequest req;
|
|
|
|
|
+ req.set_collection(config_.events_collection);
|
|
|
|
|
+ req.set_id(event_id);
|
|
|
|
|
+
|
|
|
|
|
+ auto* data = req.mutable_data();
|
|
|
|
|
+ auto* fields = data->mutable_fields();
|
|
|
|
|
+
|
|
|
|
|
+ (*fields)["status"].set_int64_value(static_cast<int64_t>(status));
|
|
|
|
|
+
|
|
|
|
|
+ if (status == ::smartbotic::llm::MESSAGE_EVENT_STATUS_COMPLETED) {
|
|
|
|
|
+ (*fields)["completed_at"].set_string_value(GetCurrentTimestamp());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!error.empty()) {
|
|
|
|
|
+ (*fields)["error"].set_string_value(error);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentResponse resp;
|
|
|
|
|
+ auto grpc_status = document_stub_->UpdateDocument(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!grpc_status.ok()) {
|
|
|
|
|
+ spdlog::warn("EventProcessor: Failed to update event status: {}",
|
|
|
|
|
+ grpc_status.error_message());
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::UpdateSessionStreamingState(const std::string& session_id,
|
|
|
|
|
+ const ::smartbotic::llm::StreamingState& state) {
|
|
|
|
|
+ // Update the session document with streaming state
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentRequest req;
|
|
|
|
|
+ req.set_collection("__llm_sessions");
|
|
|
|
|
+ req.set_id(session_id);
|
|
|
|
|
+
|
|
|
|
|
+ auto* data = req.mutable_data();
|
|
|
|
|
+ auto* fields = data->mutable_fields();
|
|
|
|
|
+
|
|
|
|
|
+ // Create streaming_state as nested map
|
|
|
|
|
+ auto* streaming_state = (*fields)["streaming_state"].mutable_map_value();
|
|
|
|
|
+ auto* ss_fields = streaming_state->mutable_fields();
|
|
|
|
|
+
|
|
|
|
|
+ (*ss_fields)["is_streaming"].set_bool_value(state.is_streaming());
|
|
|
|
|
+ (*ss_fields)["current_event_id"].set_string_value(state.current_event_id());
|
|
|
|
|
+ (*ss_fields)["last_chunk_sequence"].set_int64_value(state.last_chunk_sequence());
|
|
|
|
|
+ if (state.has_started_at()) {
|
|
|
|
|
+ (*ss_fields)["started_at"].set_string_value(TimestampToString(state.started_at()));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::UpdateDocumentResponse resp;
|
|
|
|
|
+ auto grpc_status = document_stub_->UpdateDocument(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!grpc_status.ok()) {
|
|
|
|
|
+ spdlog::warn("EventProcessor: Failed to update session streaming state: {}",
|
|
|
|
|
+ grpc_status.error_message());
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void EventProcessor::RecoverStaleEvents() {
|
|
|
|
|
+ // Find events that have been processing for too long
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::QueryRequest req;
|
|
|
|
|
+ req.set_collection(config_.events_collection);
|
|
|
|
|
+
|
|
|
|
|
+ auto* query = req.mutable_query();
|
|
|
|
|
+ auto* status_filter = query->add_filters();
|
|
|
|
|
+ status_filter->set_field("status");
|
|
|
|
|
+ status_filter->set_operator_(::smartbotic::database::FilterOperator::FILTER_OPERATOR_EQUALS);
|
|
|
|
|
+ auto* status_val = status_filter->mutable_value();
|
|
|
|
|
+ status_val->set_int64_value(
|
|
|
|
|
+ static_cast<int64_t>(::smartbotic::llm::MESSAGE_EVENT_STATUS_PROCESSING)
|
|
|
|
|
+ );
|
|
|
|
|
+
|
|
|
|
|
+ req.set_limit(100);
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::QueryResponse resp;
|
|
|
|
|
+ auto status = query_stub_->Query(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto now = std::chrono::system_clock::now();
|
|
|
|
|
+
|
|
|
|
|
+ for (const auto& doc : resp.documents()) {
|
|
|
|
|
+ auto event = DocumentToEvent(doc.data(), doc.id());
|
|
|
|
|
+
|
|
|
|
|
+ // Check if started_at is older than threshold
|
|
|
|
|
+ if (event.has_started_at()) {
|
|
|
|
|
+ auto started = std::chrono::system_clock::from_time_t(event.started_at().seconds());
|
|
|
|
|
+ auto age = std::chrono::duration_cast<std::chrono::seconds>(now - started).count();
|
|
|
|
|
+
|
|
|
|
|
+ if (age > config_.stale_event_timeout_seconds) {
|
|
|
|
|
+ spdlog::warn("EventProcessor: Recovering stale event {} (age={}s)",
|
|
|
|
|
+ event.id(), age);
|
|
|
|
|
+
|
|
|
|
|
+ // Reset to pending so it can be picked up again
|
|
|
|
|
+ UpdateEventStatus(event.id(), ::smartbotic::llm::MESSAGE_EVENT_STATUS_PENDING);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::QueueMessageEvent(const ::smartbotic::llm::MessageEvent& event)
|
|
|
|
|
+ -> EventResult<std::string> {
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::CreateDocumentRequest req;
|
|
|
|
|
+ req.set_collection(config_.events_collection);
|
|
|
|
|
+
|
|
|
|
|
+ // Set a new ID if not provided
|
|
|
|
|
+ ::smartbotic::llm::MessageEvent event_copy = event;
|
|
|
|
|
+ if (event_copy.id().empty()) {
|
|
|
|
|
+ event_copy.set_id(GenerateId());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ *req.mutable_data() = EventToDocument(event_copy);
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::CreateDocumentResponse resp;
|
|
|
|
|
+ auto status = document_stub_->CreateDocument(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ return EventResult<std::string>::Error(status.error_message());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return EventResult<std::string>::Ok(event_copy.id());
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::GetStreamingState(const std::string& session_id)
|
|
|
|
|
+ -> EventResult<::smartbotic::llm::StreamingState> {
|
|
|
|
|
+
|
|
|
|
|
+ auto result = session_store_->GetSession(session_id, "", false);
|
|
|
|
|
+ if (!result.ok) {
|
|
|
|
|
+ return EventResult<::smartbotic::llm::StreamingState>::Error(result.error);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Extract streaming_state from session (would need to be added to Session proto)
|
|
|
|
|
+ ::smartbotic::llm::StreamingState state;
|
|
|
|
|
+ // The streaming_state would be part of the session document
|
|
|
|
|
+ return EventResult<::smartbotic::llm::StreamingState>::Ok(state);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::SyncChunks(const std::string& session_id,
|
|
|
|
|
+ const std::string& event_id,
|
|
|
|
|
+ int32_t since_sequence)
|
|
|
|
|
+ -> EventResult<std::vector<::smartbotic::llm::ResponseChunk>> {
|
|
|
|
|
+
|
|
|
|
|
+ std::vector<::smartbotic::llm::ResponseChunk> chunks;
|
|
|
|
|
+
|
|
|
|
|
+ grpc::ClientContext ctx;
|
|
|
|
|
+ ::smartbotic::database::QueryRequest req;
|
|
|
|
|
+ req.set_collection(config_.chunks_collection);
|
|
|
|
|
+
|
|
|
|
|
+ auto* query = req.mutable_query();
|
|
|
|
|
+
|
|
|
|
|
+ // Filter by event_id
|
|
|
|
|
+ auto* event_filter = query->add_filters();
|
|
|
|
|
+ event_filter->set_field("message_event_id");
|
|
|
|
|
+ event_filter->set_operator_(::smartbotic::database::FilterOperator::FILTER_OPERATOR_EQUALS);
|
|
|
|
|
+ event_filter->mutable_value()->set_string_value(event_id);
|
|
|
|
|
+
|
|
|
|
|
+ // Filter by sequence > since_sequence
|
|
|
|
|
+ auto* seq_filter = query->add_filters();
|
|
|
|
|
+ seq_filter->set_field("sequence");
|
|
|
|
|
+ seq_filter->set_operator_(::smartbotic::database::FilterOperator::FILTER_OPERATOR_GREATER_THAN);
|
|
|
|
|
+ seq_filter->mutable_value()->set_int64_value(since_sequence);
|
|
|
|
|
+
|
|
|
|
|
+ // Order by sequence
|
|
|
|
|
+ auto* order = query->add_order_by();
|
|
|
|
|
+ order->set_field("sequence");
|
|
|
|
|
+ order->set_direction(::smartbotic::database::SortDirection::SORT_DIRECTION_ASCENDING);
|
|
|
|
|
+
|
|
|
|
|
+ req.set_limit(1000);
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::QueryResponse resp;
|
|
|
|
|
+ auto status = query_stub_->Query(&ctx, req, &resp);
|
|
|
|
|
+
|
|
|
|
|
+ if (!status.ok()) {
|
|
|
|
|
+ return EventResult<std::vector<::smartbotic::llm::ResponseChunk>>::Error(
|
|
|
|
|
+ status.error_message()
|
|
|
|
|
+ );
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ for (const auto& doc : resp.documents()) {
|
|
|
|
|
+ // Convert document to ResponseChunk
|
|
|
|
|
+ // This would require a DocumentToChunk method
|
|
|
|
|
+ ::smartbotic::llm::ResponseChunk chunk;
|
|
|
|
|
+ // Parse fields from doc.data()
|
|
|
|
|
+ chunks.push_back(chunk);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return EventResult<std::vector<::smartbotic::llm::ResponseChunk>>::Ok(std::move(chunks));
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::EventToDocument(const ::smartbotic::llm::MessageEvent& event)
|
|
|
|
|
+ -> ::smartbotic::database::MapValue {
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::MapValue data;
|
|
|
|
|
+ auto* fields = data.mutable_fields();
|
|
|
|
|
+
|
|
|
|
|
+ (*fields)["id"].set_string_value(event.id());
|
|
|
|
|
+ (*fields)["workspace_id"].set_string_value(event.workspace_id());
|
|
|
|
|
+ (*fields)["session_id"].set_string_value(event.session_id());
|
|
|
|
|
+ (*fields)["user_id"].set_string_value(event.user_id());
|
|
|
|
|
+ (*fields)["event_type"].set_int64_value(static_cast<int64_t>(event.event_type()));
|
|
|
|
|
+ (*fields)["status"].set_int64_value(static_cast<int64_t>(event.status()));
|
|
|
|
|
+
|
|
|
|
|
+ // Request parameters
|
|
|
|
|
+ (*fields)["model_id"].set_string_value(event.model_id());
|
|
|
|
|
+ (*fields)["provider_id"].set_string_value(event.provider_id());
|
|
|
|
|
+ (*fields)["temperature"].set_double_value(event.temperature());
|
|
|
|
|
+ (*fields)["max_tokens"].set_int64_value(event.max_tokens());
|
|
|
|
|
+
|
|
|
|
|
+ // Processing metadata
|
|
|
|
|
+ (*fields)["processing_node_id"].set_string_value(event.processing_node_id());
|
|
|
|
|
+ (*fields)["created_at"].set_string_value(GetCurrentTimestamp());
|
|
|
|
|
+
|
|
|
|
|
+ if (!event.error().empty()) {
|
|
|
|
|
+ (*fields)["error"].set_string_value(event.error());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Store user_message as nested map if present
|
|
|
|
|
+ if (event.has_user_message()) {
|
|
|
|
|
+ // Serialize message (simplified)
|
|
|
|
|
+ auto* msg_map = (*fields)["user_message"].mutable_map_value();
|
|
|
|
|
+ auto* msg_fields = msg_map->mutable_fields();
|
|
|
|
|
+ (*msg_fields)["id"].set_string_value(event.user_message().id());
|
|
|
|
|
+ (*msg_fields)["role"].set_int64_value(
|
|
|
|
|
+ static_cast<int64_t>(event.user_message().role())
|
|
|
|
|
+ );
|
|
|
|
|
+ // Add content array
|
|
|
|
|
+ auto* content_arr = (*msg_fields)["content"].mutable_array_value();
|
|
|
|
|
+ for (const auto& part : event.user_message().content()) {
|
|
|
|
|
+ auto* part_val = content_arr->add_values();
|
|
|
|
|
+ auto* part_map = part_val->mutable_map_value();
|
|
|
|
|
+ auto* part_fields = part_map->mutable_fields();
|
|
|
|
|
+ if (part.has_text()) {
|
|
|
|
|
+ (*part_fields)["type"].set_string_value("text");
|
|
|
|
|
+ (*part_fields)["text"].set_string_value(part.text());
|
|
|
|
|
+ }
|
|
|
|
|
+ // Handle other content types
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Store tool_results as array if present
|
|
|
|
|
+ if (event.tool_results_size() > 0) {
|
|
|
|
|
+ auto* results_arr = (*fields)["tool_results"].mutable_array_value();
|
|
|
|
|
+ for (const auto& result : event.tool_results()) {
|
|
|
|
|
+ auto* result_val = results_arr->add_values();
|
|
|
|
|
+ auto* result_map = result_val->mutable_map_value();
|
|
|
|
|
+ auto* result_fields = result_map->mutable_fields();
|
|
|
|
|
+ (*result_fields)["tool_call_id"].set_string_value(result.tool_call_id());
|
|
|
|
|
+ (*result_fields)["content"].set_string_value(result.content());
|
|
|
|
|
+ (*result_fields)["is_error"].set_bool_value(result.is_error());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Store tools as array
|
|
|
|
|
+ if (event.tools_size() > 0) {
|
|
|
|
|
+ auto* tools_arr = (*fields)["tools"].mutable_array_value();
|
|
|
|
|
+ for (const auto& tool : event.tools()) {
|
|
|
|
|
+ auto* tool_val = tools_arr->add_values();
|
|
|
|
|
+ auto* tool_map = tool_val->mutable_map_value();
|
|
|
|
|
+ auto* tool_fields = tool_map->mutable_fields();
|
|
|
|
|
+ (*tool_fields)["name"].set_string_value(tool.name());
|
|
|
|
|
+ (*tool_fields)["description"].set_string_value(tool.description());
|
|
|
|
|
+ (*tool_fields)["input_schema"].set_string_value(tool.input_schema());
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return data;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::DocumentToEvent(const ::smartbotic::database::MapValue& data,
|
|
|
|
|
+ const std::string& doc_id)
|
|
|
|
|
+ -> ::smartbotic::llm::MessageEvent {
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::llm::MessageEvent event;
|
|
|
|
|
+ event.set_id(doc_id);
|
|
|
|
|
+
|
|
|
|
|
+ const auto& fields = data.fields();
|
|
|
|
|
+
|
|
|
|
|
+ if (auto it = fields.find("workspace_id"); it != fields.end()) {
|
|
|
|
|
+ event.set_workspace_id(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("session_id"); it != fields.end()) {
|
|
|
|
|
+ event.set_session_id(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("user_id"); it != fields.end()) {
|
|
|
|
|
+ event.set_user_id(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("event_type"); it != fields.end()) {
|
|
|
|
|
+ event.set_event_type(
|
|
|
|
|
+ static_cast<::smartbotic::llm::MessageEventType>(it->second.int64_value())
|
|
|
|
|
+ );
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("status"); it != fields.end()) {
|
|
|
|
|
+ event.set_status(
|
|
|
|
|
+ static_cast<::smartbotic::llm::MessageEventStatus>(it->second.int64_value())
|
|
|
|
|
+ );
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("model_id"); it != fields.end()) {
|
|
|
|
|
+ event.set_model_id(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("provider_id"); it != fields.end()) {
|
|
|
|
|
+ event.set_provider_id(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("temperature"); it != fields.end()) {
|
|
|
|
|
+ event.set_temperature(static_cast<float>(it->second.double_value()));
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("max_tokens"); it != fields.end()) {
|
|
|
|
|
+ event.set_max_tokens(static_cast<int32_t>(it->second.int64_value()));
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("processing_node_id"); it != fields.end()) {
|
|
|
|
|
+ event.set_processing_node_id(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("error"); it != fields.end()) {
|
|
|
|
|
+ event.set_error(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Parse timestamps
|
|
|
|
|
+ if (auto it = fields.find("created_at"); it != fields.end()) {
|
|
|
|
|
+ *event.mutable_created_at() = StringToTimestamp(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("started_at"); it != fields.end()) {
|
|
|
|
|
+ *event.mutable_started_at() = StringToTimestamp(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (auto it = fields.find("completed_at"); it != fields.end()) {
|
|
|
|
|
+ *event.mutable_completed_at() = StringToTimestamp(it->second.string_value());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Parse user_message (simplified)
|
|
|
|
|
+ if (auto it = fields.find("user_message"); it != fields.end() &&
|
|
|
|
|
+ it->second.has_map_value()) {
|
|
|
|
|
+ // Would need full message parsing
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Parse tool_results (simplified)
|
|
|
|
|
+ if (auto it = fields.find("tool_results"); it != fields.end() &&
|
|
|
|
|
+ it->second.has_array_value()) {
|
|
|
|
|
+ // Would need full array parsing
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return event;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+auto EventProcessor::ChunkToDocument(const ::smartbotic::llm::ResponseChunk& chunk)
|
|
|
|
|
+ -> ::smartbotic::database::MapValue {
|
|
|
|
|
+
|
|
|
|
|
+ ::smartbotic::database::MapValue data;
|
|
|
|
|
+ auto* fields = data.mutable_fields();
|
|
|
|
|
+
|
|
|
|
|
+ (*fields)["id"].set_string_value(chunk.id());
|
|
|
|
|
+ (*fields)["workspace_id"].set_string_value(chunk.workspace_id());
|
|
|
|
|
+ (*fields)["session_id"].set_string_value(chunk.session_id());
|
|
|
|
|
+ (*fields)["message_event_id"].set_string_value(chunk.message_event_id());
|
|
|
|
|
+ (*fields)["message_id"].set_string_value(chunk.message_id());
|
|
|
|
|
+ (*fields)["sequence"].set_int64_value(chunk.sequence());
|
|
|
|
|
+ (*fields)["chunk_type"].set_int64_value(static_cast<int64_t>(chunk.chunk_type()));
|
|
|
|
|
+ (*fields)["is_final"].set_bool_value(chunk.is_final());
|
|
|
|
|
+ (*fields)["created_at"].set_string_value(
|
|
|
|
|
+ chunk.has_created_at() ? TimestampToString(chunk.created_at()) : GetCurrentTimestamp()
|
|
|
|
|
+ );
|
|
|
|
|
+
|
|
|
|
|
+ // Content based on type
|
|
|
|
|
+ if (!chunk.content_delta().empty()) {
|
|
|
|
|
+ (*fields)["content_delta"].set_string_value(chunk.content_delta());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (!chunk.thinking_delta().empty()) {
|
|
|
|
|
+ (*fields)["thinking_delta"].set_string_value(chunk.thinking_delta());
|
|
|
|
|
+ }
|
|
|
|
|
+ if (chunk.has_tool_call()) {
|
|
|
|
|
+ auto* tc_map = (*fields)["tool_call"].mutable_map_value();
|
|
|
|
|
+ auto* tc_fields = tc_map->mutable_fields();
|
|
|
|
|
+ (*tc_fields)["id"].set_string_value(chunk.tool_call().id());
|
|
|
|
|
+ (*tc_fields)["name"].set_string_value(chunk.tool_call().name());
|
|
|
|
|
+ (*tc_fields)["arguments"].set_string_value(chunk.tool_call().arguments());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return data;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+} // namespace smartbotic::llm::event
|