Skip to content

Latest commit

 

History

History
189 lines (153 loc) · 14.2 KB

File metadata and controls

189 lines (153 loc) · 14.2 KB

Wasman: Durable WebAssembly (WASM) Execution Engine

Надежный, легко переиспользуемый и производительный движок устойчивого выполнения (Durable Execution) на языке Go с использованием WebAssembly (WASM) и среды выполнения wazero. Он обеспечивает отказоустойчивое исполнение бизнес-логики в изолированной песочнице с автоматическим сохранением снимков памяти (snapshotting), прозрачным восстановлением после сбоев и эффективным по памяти потоковым обменом данными, полностью избавляя от зависимостей CGO и glibc.


Философия Durable Execution (Устойчивого выполнения)

Современные распределенные архитектуры часто требуют выполнения длительных или многошаговых бизнес-процессов (интеграций, цепочек вызовов, оркестраций), которые должны переживать сбои инфраструктуры. Традиционные подходы полагаются на сложные машины состояний в базах данных или тяжелые внешние оркестраторы.

Wasman решает эту проблему за счет изолированной линейной памяти WebAssembly:

  1. Изоляция хоста и гостя (Host-Guest Isolation): Бизнес-логика компилируется в независимый .wasm модуль (на TinyGo/Go) и выполняется внутри чистого Go-рантайма wazero.
  2. Stateless-хост и надежное хранилище: Узлы, выполняющие код, остаются полностью stateless. Все состояние процесса (снимки памяти, журнал вызовов oplog) сохраняется во внешнем S3-совместимом или файловом хранилище.
  3. Принцип «Черного ящика» (Black Box): Разработчики, использующие платформу (например, NativeBPM), работают только со сгенерированными API-клиентами (Protobuf/GRPC или SDK). Вся сложность WebAssembly, создания снимков памяти и транзакционного контроля полностью скрыта от них.
       [ Запрос клиента / API ] (StartProcess / CompleteTask)
                 
                 
       ┌──────────────────┐
         WorkflowEngine   (Хост-оркестратор)
       └─────────┬────────┘
                 
      ┌──────────┴──────────┐
                           
┌───────────┐         ┌───────────┐
 bpmn_vm    (WASM)    worker    (WASM бизнес-логика)
 Интерпретатор         Воркер   
└─────┬─────┘         └─────┬─────┘
                           
      └──────────┬──────────┘
                  (Точки сохранения / Чекпоинты)
                 
       ┌──────────────────┐
         Snapshot Store   (Сжатые снимки и дельты памяти)
       └─────────┬────────┘
                 
        ┌────────┴────────┐
                         
   [ Хранилище S3 ]  [ Локальные файлы ]

Архитектурные особенности и технологии

1. Обязательное сжатие данных (Gzip)

Создание снимков памяти WASM-модулей генерирует дампы линейной памяти (кратно страницам по 64 КБ). Для оптимизации дискового пространства и трафика:

  • Gzip-сжатие: Все снимки памяти, постраничные дельты изменений и логи oplog автоматически сжимаются на лету с использованием gzip.
  • Строгий формат: Чтение данных требует обязательного наличия gzip-сжатия. Несжатые сырые снимки не поддерживаются, что гарантирует максимальную эффективность хранения для всех состояний процессов.

2. Стриминг в памяти с $O(1)$ RAM

При передаче больших объемов данных (файлов, тяжелых JSON/CSV):

  • Данные записываются и считываются из памяти WASM блоками через буферы потоков.
  • Это гарантирует фиксированное и минимальное потребление RAM ($O(1)$) вне зависимости от размера обрабатываемых файлов, исключая перегрузку сборщика мусора (GC) и утечки памяти.
  • Все взаимодействия происходят полностью в оперативной памяти через переданные Go-обработчики (handlers), полностью избавляя от необходимости локальных HTTP-вызовов и открытия сетевых портов.

3. Постраничные дельты памяти (Delta Snapshots)

Вместо записи полного многомегабайтного снимка памяти при каждом сохранении:

  • Wasman хэширует каждую страницу памяти (64 КБ) по алгоритму FNV-64a.
  • На чекпоинтах сохраняются только те страницы, хэш которых изменился (грязные страницы), что на порядок снижает нагрузку на диск и сеть.

4. Оптимистичная блокировка (OCC)

В конкурентной распределенной среде, когда несколько узлов могут одновременно получить триггер выполнения для одного и того же экземпляра процесса:

  • S3-клиент использует механизм HTTP-заголовков If-Match на основе ETag объектов.
  • Если другое приложение успело обновить снимок процесса раньше, текущая запись падает с ошибкой OCC, предотвращая порчу данных и гонки состояний.

Обработка сбоев в Corner Cases

Wasman гарантирует надежное восстановление состояния гостевой системы после аварийных завершений:

Сценарий: Падение сервера во время выполнения процесса

  1. До шага: Воркер начинает выполнять процесс. Он доходит до точки сохранения (чекпоинта) — например, перед отправкой внешнего HTTP-запроса или переходом в режим ожидания задачи пользователя (User Task).
  2. Чекпоинт:
    • Движок приостанавливает выполнение.
    • Делается снимок состояния памяти (Full Snapshot или Delta Snapshot), который пишется в хранилище.
    • Записывается транзакционный лог перехода в oplog.
  3. Сбой: Происходит аварийная остановка хост-сервера (сбой по питанию, OOM или плановый перезапуск контейнера).
  4. Восстановление:
    • Новый хост получает команду продолжить процесс.
    • Он считывает метаданные из S3, загружает исходный WASM-модуль и подтягивает сжатый снимок памяти.
    • Линейная память WASM VM мгновенно восстанавливается до состояния на момент последнего чекпоинта.
    • Проигрываются логи oplog для восстановления промежуточных вызовов, и выполнение продолжается бесшовно и незаметно для внешней системы.

Структура каталогов

  • wasman.go: Инициализация, компиляция и основной цикл выполнения WASM.
  • compress.go: Вспомогательные утилиты для прозрачного gzip-сжатия.
  • fs_store.go: Файловое хранилище снимков с поддержкой сжатия.
  • s3_store.go: Хранилище снимков в S3 с поддержкой сжатия и OCC.
  • types.go: Базовые интерфейсы, структуры и ошибки.
  • examples/: Примеры интеграций:
    • process-csv/: Стриминг и парсинг CSV с эмуляцией сбоя и восстановлением.
    • camunda/: Воркер для интеграции с внешними задачами Camunda 7.
    • temporal/: CRM/CRM-операции с чекпоинтами в Temporal.
    • gotenberg-telegram/: Потоковый бот конвертации файлов в PDF.
    • s3-store/: Демонстрация сохранения снимков напрямую в S3/MinIO.
    • in-memory-channel/: Исключительно внутрипамятый потоковый обмен данными между хостом и WASM-гостем, полностью исключающий сетевые вызовы по TCP loopback.
    • safe-task/: Выполнение изолированных задач с использованием высокоуровневой и безопасной утилиты-раннера RunTask.
    • wasm-inspector/: Низкоуровневая утилита для инспекции и запуска гостевых WASM-модулей с настраиваемыми параметрами WASI.

Быстрый старт

Запуск тестов

Для запуска тестов ядра:

go test -v .

Запуск примера со сбоем и восстановлением (CSV)

Пример process-csv наглядно демонстрирует цикл падения и восстановления процесса:

  1. Скомпилируйте WASM-воркер:
    make build-worker
  2. Запустите пример:
    make run-csv-example

В ходе выполнения:

  • Запустится локальный мок-сервер.
  • Начнется чтение CSV. На первом чекпоинте сработает эмуляция сбоя (WithCrash(true)).
  • Программа завершится, а на диске сохранится сжатый снимок памяти.
  • При повторном запуске движок восстановит состояние из снимка и успешно завершит обработку CSV-файла.

Пример кода интеграции

package main

import (
	"context"
	"fmt"
	"github.com/nativebpm/wasman"
)

func main() {
	// 1. Инициализируем хранилище со включенным сжатием
	store := &wasman.FileSnapshotStore{
		Dir:         "snapshots",
		Compression: true,
	}

	// 2. Определяем обработчики стриминга в памяти
	downloadHandler := func() ([]byte, error) {
		return []byte("my input data stream"), nil
	}
	uploadHandler := func(payload []byte) error {
		fmt.Printf("Получены выходные данные: %s\n", string(payload))
		return nil
	}

	// 3. Запускаем сессию с использованием высокоуровневого Fluent Runner API.
	// Если снимок для "my-session-id" существует, память восстановится автоматически.
	crashed, err := wasman.NewRunner().
		WithWasmPath("worker.wasm").
		WithStore(store).
		WithSessionID("my-session-id").
		WithEntrypoint("run").
		WithDownloadHandler(downloadHandler).
		WithUploadHandler(uploadHandler).
		Run()

	if err != nil {
		if crashed {
			fmt.Println("Выполнение приостановлено на контрольной точке.")
		} else {
			fmt.Printf("Ошибка выполнения: %v\n", err)
		}
	} else {
		fmt.Println("Выполнение успешно завершено!")
	}
}

Производительность и бенчмарки

Подробные профили тестирования производительности CPU и памяти (с упором на оптимизированное горячее выполнение JIT-модулей со скоростью около ~38 мкс) задокументированы в файле Benchmarks & Profiling Profile.