Итак в прошлый раз я рассказал вам про свой проект в целом ,сегодня же мы будем разбирать бекенд
Глубокий разбор архитектуры Zettaverse AgentHub
-
Кастомная валидация структуры flow_json и построение DAG
Архитектура DAG-оркестрации
В основе AgentHub лежит Directed Acyclic Graph (DAG) — направленный ациклический граф, где узлы представляют собой задачи агентов, а рёбра — зависимости между ними . Это позволяет декомпозировать сложные задачи на подзадачи, назначать каждой специализированного агента и параллельно выполнять независимые ветви . Это схема задач.( Схему задач делал в конструкторе )
┌─────────────┐ │ Task A │ (root) └──────┬──────┘ │ ┌────────────┼────────────┐ ▼ ▼ ▼ ┌───────────┐ ┌───────────┐ ┌───────────┐ │ Task B │ │ Task C │ │ Task D │ (параллельная группа) └─────┬─────┘ └─────┬─────┘ └─────┬─────┘ │ │ │ └─────────────┼─────────────┘ ▼ ┌─────────────┐ │ Merge │ (terminal) └─────────────┘
Структура flow_json
Формат определения workflow выглядит так :
json { “name”: “market-analysis”, “description”: “Comprehensive market analysis for product launch”, “nodes”: [ { “id”: “research”, “task”: “Research competitor landscape and market size”, “inputs”: [“product_description”], “outputs”: [“competitor_list”, “market_size”] }, { “id”: “pricing”, “task”: “Collect pricing and feature data from competitors”, “inputs”: [“competitor_list”], “outputs”: [“pricing_data”, “feature_matrix”] }, { “id”: “strategy”, “task”: “Analyze positioning opportunities and pricing strategy”, “inputs”: [“pricing_data”, “feature_matrix”, “market_size”], “outputs”: [“positioning_report”, “pricing_recommendation”] }, { “id”: “merge”, “task”: “Write executive summary combining all findings”, “inputs”: [“positioning_report”, “pricing_recommendation”], “outputs”: [“executive_summary”] } ], “max_parallel”: 5, “timeout_per_agent”: 300, “retry_on_failure”: true, “quality_threshold”: 0.7 }
Процесс валидации и построения графа
-
Синтаксическая валидация — проверка корректности JSON-схемы
-
Семантическая валидация : — Отсутствие циклов в графе зависимостей — Все входные данные производятся вышестоящими узлами — Нет недостижимых узлов — Критический путь в допустимых пределах
-
Топологическая сортировка — определение порядка выполнения
-
Выявление параллельных групп — узлы без взаимозависимостей выполняются одновременно
Состояния агентов
Состояние Описание (PENDING Ожидание зависимостей) (READY Все зависимости выполнены, в очереди) (RUNNING В процессе выполнения) (COMPLETED Успешно завершён) (FAILED Провален после всех попыток) (EVALUATING Идёт оценка качества вывода)
-
CGO и управление памятью: хардкорный разбор
Архитектура взаимодействия
AgentHub использует Rust MCP engine как cdylib для CGO , что обеспечивает высокопроизводительную обработку JSON-RPC 2.0 запросов.
Критические моменты управления памятью
-
Преобразование строк:
C.CString
package main
/* #include <stdlib.h> typedef struct { char* data; int len; } RustString;
void free_rust_string(RustString* s) { if (s != NULL && s->data != NULL) { free(s->data); } } */ import “C” import “unsafe”
// Вызов Rust-функции с передачей строки func callRustFunction(input string) (string, error) { // 1. Выделяем память в C-стиле (не управляется GC!) cInput := C.CString(input) // 2. Гарантируем освобождение при выходе из функции defer C.free(unsafe.Pointer(cInput)) // 3. Вызов Rust-функции var result C.RustString ret := C.rust_process_string(cInput, &result) // 4. Освобождаем память, выделенную Rust-стороной defer C.free_rust_string(&result) // 5. Конвертируем результат в Go-строку (копирование) goResult := C.GoStringN(result.data, C.int(result.len)) return goResult, nil }
-
Почему
deferкритичен
Использование defer гарантирует освобождение памяти даже при панике:
go func riskyOperation() { cStr := C.CString(“large data”) // Выделение defer C.free(unsafe.Pointer(cStr)) // Гарантированное освобождение // Даже если здесь panic, память будет освобождена doSomethingDangerous(cStr) }
-
Паттерн
free_rust_string
Rust и Go используют разные аллокаторы. Rust не может напрямую освобождать память, выделенную C.malloc, поэтому нужен bridge :
c // В C-обёртке void free_rust_string(RustString* s) { if (s != NULL && s->data != NULL) { free(s->data); // Используем тот же аллокатор, что при выделении s->data = NULL; s->len = 0; } }
-
Доказательство отсутствия утечек
go // Мониторинг памяти для отладки import “runtime”
func debugMemory() { var m runtime.MemStats runtime.ReadMemStats(&m) fmt.Printf(“Alloc = %v MiB”, m.Alloc / 1024 / 1024) }
// Использование в тестах func TestNoMemoryLeaks(t *testing.T) { for i := 0; i < 10000; i++ { callRustFunction(«test string » + strconv.Itoa(i)) } debugMemory() // Должно оставаться стабильным }
-
Реальный боевой кейс: управление агентами через WebSockets
Сценарий
Один управляющий агент (через Vue 3 интерфейс) оркестрирует других агентов в реальном времени .
Архитектура взаимодействия
┌─────────────────────────────────────────────────────────────┐ │ Vue 3 Frontend │ │ ┌─────────────┐ ┌─────────────┐ ┌──────────────────┐ │ │ │ Dashboard │ │ Flow Editor │ │ Live Chat │ │ │ └──────┬──────┘ └──────┬──────┘ └────────┬─────────┘ │ └─────────┼────────────────┼───────────────────┼────────────┘ │ │ │ └────────────────┼───────────────────┘ │ WebSocket (WSS) ▼ ┌─────────────────────────────────────────────────────────────┐ │ agent-hub-core (Go Backend) │ │ ┌──────────────────────────────────────────────────────┐ │ │ │ WebSocket Hub / Orchestrator │ │ │ │ ┌──────┐ ┌──────┐ ┌──────┐ ┌─────────────────┐│ │ │ │ │Agent1│ │Agent2│ │Agent3│ │ MCP Engine ││ │ │ │ └──┬───┘ └──┬───┘ └──┬───┘ │ (Rust via CGO) ││ │ │ └─────┼─────────┼─────────┼──────┴─────────────────┘│ │ └────────┼─────────┼─────────┼──────────────────────────┘ │ │ │ ▼ ▼ ▼ ┌────────┐┌────────┐┌────────┐ │Agent A ││Agent B ││Agent C │ (Удалённые агенты) └────────┘└────────┘└────────┘
Поток событий
-
Пользователь через Vue 3 интерфейс создаёт workflow
-
Управляющий агент (в Go backend) транслирует задачу через WebSocket
-
Агенты-исполнители получают задачу и начинают работу
-
Статусы обновляются в реальном времени через WebSocket
Типы WebSocket событий
typescript // Из enhanced-chat implementation type WebSocketEvent = | { type: ‘chat’, data: ChatEvent } // OpenClaw gateway events | { type: ‘agent’, data: AgentStream } // Agent stream events | { type: ‘activity’, data: Activity } // Status updates | { type: ‘tool_call’, data: ToolCall } // Tool execution
Жизненный цикл выполнения
Init → Run → Spawn (parallel) → Eval → Merge → Status
-
Init — определение DAG workflow
-
Run — запуск оркестрации
-
Spawn — создание агентов (параллельно, до
max_parallel) -
Eval — проверка качества вывода
-
Merge — объединение результатов
-
Status — отчёт о состоянии выполнения
Компоненты Vue 3
-
Dashboard — мониторинг состояния агентов
-
Flow Editor — визуальное редактирование DAG
-
Live Chat — реальное время общения с агентами
Заключение
Zettaverse AgentHub — это многоуровневая система, где:
-
Go backend (
agent-hub-core) управляет DAG-оркестрацией с валидациейflow_json -
Rust MCP engine (
agent-hub-mcp) через CGO обеспечивает производительную обработку [citation:agent-hub-mcp] -
Vue 3 UI (
agent-hub-ui) предоставляет интерфейс реального времени через WebSockets
Ключевая фитча — превращение GitHub-подхода из человеко-ориентированного в агент-ориентированный, с DAG вместо линейной истории и сообщениями как основным средством координации . Для тех кто не в теме вот репозиторий https://github.com/orgs/Zettaverse
ссылка на оригинал статьи https://habr.com/ru/articles/1078278/