SDK и жизненный цикл
Handshake, чтение, публикация, ограниченные повторы и остановка.
Подключение
Фрагменты ниже показывают реальные методы SDK, но не заменяют полный run() из src/main.rs. SDK сверяет ID, версию и Plugin API. session.start содержит конфигурацию и фактические grants экземпляра.
let session = PluginSession::connect(PLUGIN_ID, env!("CARGO_PKG_VERSION")).await?;
let config: Config = session.start.config()?;Готовность
Проверьте config и оба точных grant до client.ready(). SDK сам посылает heartbeat каждые 10 секунд: отдельный heartbeat-поток не нужен. Ошибки startup, handshake и config фатальны; работать без этих условий нельзя.
Разовое чтение и публикация
resolve_usepi остается отдельным API разового чтения и требует usepi.read grant. Текущий manifest стартера выдает только usepi.subscribe: этот фрагмент показывает альтернативу, а полный пример работает по событиям ниже. resolve_usepi возвращает Vec<u8> с JSON.
Стартер принимает JSON-число 21.5 и строку "21.5"; окружающие пробелы внутри строки удаляются. Объект {"value":21.5}, null, bool, массив, NaN, бесконечность, шестнадцатеричная запись и десятичная запятая отклоняются, а не превращаются в ноль. При scale = 2 и offset = -1 из 21.5 получится {"value":42.0}. Вход ограничен 4096 байтами, выход 256 байтами, модуль конечного результата 1e12. f64 не подходит для точных денежных расчетов; пример не является safety-контроллером.
let raw = client
.resolve_usepi(&config.source_server_id, config.source_path.clone())
.await?;
let payload = transform(&raw, config)?;
client
.publish_mqtt(
&config.destination_server_id,
&config.destination_topic,
payload,
false,
)
.await?;Подписка Plugin API 1.1
UsepiEventKind экспортируется SDK. Подписка принимает только точный лист и отдельный usepi.subscribe grant. Первое событие - Snapshot; found=false означает отсутствие листа и отличается от JSON null. Следом приходят Update. Reset, Lagged и Cancelled завершают поток. Revisions растут по всему дереву: пропуски допустимы. Это текущее состояние, не полный журнал всех измерений.
Отмена ожидания next() сохраняет один незавершенный запрос внутри stream. Следующий next() продолжает его. Для подтвержденной отмены вызовите stream.cancel().await; отмененный cancel() также продолжается следующим cancel() или next(), без доставки новых значений. Drop посылает только best-effort cancel; при переполнении очереди освобождение обеспечивает lease хоста. Отмена subscribe_usepi до получения ID тоже может оставить подписку до истечения lease. Полный стартер обрабатывает ошибки и повторное открытие с backoff.
let mut stream = client
.subscribe_usepi(&config.source_server_id, config.source_path.clone())
.await?;
while let Some(event) = stream.next().await? {
match UsepiEventKind::try_from(event.kind) {
Ok(UsepiEventKind::Snapshot | UsepiEventKind::Update) if event.found => {
let payload = transform(&event.value_json, config)?;
client.publish_mqtt(
&config.destination_server_id,
&config.destination_topic,
payload,
false,
).await?;
}
Ok(UsepiEventKind::Reset | UsepiEventKind::Lagged | UsepiEventKind::Cancelled) => break,
_ => {}
}
}Семантика доставки
Успех publish_mqtt означает принятие запроса Host API. Используется QoS 0, false означает retain = false. Это не подтверждение доставки подписчику, не exactly-once и не запись обратно в USEPI. Публикация сама по себе не создает модель или виджет.
Подписка передает snapshot и изменения, но не является постоянным журналом или гарантией возраста измерения. После reset/lagged стартер открывает новую подписку и получает новый snapshot. Для контроля свежести нужен отдельный контракт качества и времени источника.
Цикл и повторы
Стартер обрабатывает события последовательно, без polling-таймера и неограниченной очереди. intervalMs задает основу backoff повторного открытия подписки: минимум 5 секунд, удвоение до 60 секунд. Ошибка преобразования или неясный результат MQTT не вызывают немедленного повтора старой публикации: стартер ждет следующего события.
Не оборачивайте SDK-вызовы в короткий timeout с немедленным повтором: отмененный future не гарантирует отмену принятой операции. Долгие CPU-вычисления и блокирующий I/O задерживают heartbeat и остановку.
Остановка
wait_for_shutdown() слушается во внешнем tokio::select! одновременно со всей работой, включая ready, чтение, публикацию и паузу. При shutdown или закрытии канала ожидания отменяются, журнал получает STOP, процесс выходит. Уже принятая MQTT-публикация может завершиться после остановки: отмена ожидания ее не откатывает.
Config приходит при запуске, не как произвольный hot reload. Для учебного опыта сначала остановите экземпляр, затем сохраните изменения и запустите снова.
Пределы примера и Host API
| Ресурс | Community Starter | Текущий потолок Host API |
|---|---|---|
| Входной JSON | 4096 байт | USEPI value: 192 KiB |
| Выходной JSON | 256 байт | MQTT payload: 192 KiB |
| Результат | Конечное число, модуль до 1e12 | Зависит от операции |
| Путь подписки | 1..16 сегментов | До 64 сегментов; весь usepi.subscribe scope до 256 байт |
| Топик | 1..256 символов | 1024 байта |
Ресурсы и журнал
Текущие defaults процесса: мягкая память 192 MiB, максимум 256 MiB, CPU 100%, до 64 tasks; политика установки может быть строже. Ноль в редакторе лимитов означает default, не безлимит. Стандартный deadline остановки - 10 секунд, но реагировать следует сразу.
Предупреждения стартера ограничены одним сообщением в минуту. В журнал попадают этап и безопасный код, не config, payload или произвольная ошибка Host API. Не добавляйте бесконечные коллекции или очереди и не рассчитывайте на API квот постоянного хранилища, которого у стартера нет.