变更记录
- 11/23/2020:初始草案
- 10/06/2022:引入基于
hashicorp/go-plugin的插件系统 - 10/14/2022:
- 添加
ListenCommit,将一个区块中的状态写入压平成单个批次。 - 从缓存存储中移除监听器,只应监听
rootmulti.Store。 - 移除
HaltAppOnDeliveryError(),错误默认会向上传播;如果实现不希望传播错误,应返回 nil。
- 添加
- 26/05/2023:更新为 ABCI 2.0
状态
提议中摘要
本 ADR 定义了一组变更,用于支持监听各个 KVStore 的状态变化,并将这些数据暴露给消费者。背景
当前,KVStore 数据可以通过查询进行远程访问, 这些查询要么通过 Tendermint 和 ABCI 处理,要么通过 gRPC 服务器处理。 除了这种请求/响应式查询之外,如果还能有一种方式在状态变化发生时实时监听这些变化,将会很有价值。决策
我们将修改CommitMultiStore 接口及其具体的 rootmulti 实现,并引入新的 listenkv.Store,以支持监听底层 KVStore 中的状态变化。我们不需要监听缓存存储,因为无法确定这些写入最终一定会被提交,而且这些写入最终也会在 rootmulti.Store 中重复出现,因此我们只应监听 rootmulti.Store。
我们将引入一个插件系统,用于配置和运行流式服务,把这些状态变化及其周边的 ABCI 消息上下文写入不同的目标位置。
监听
在新文件store/types/listening.go 中,我们将创建一个 MemoryListener 结构体,用于从 KVStore 中流式输出采用 protobuf 编码的 KV 键值对状态变化。
MemoryListener 将由具体的 rootmulti 实现在内部使用,用于收集来自 KVStore 的状态变化。
ListenKVStore
我们将创建一个新的Store 类型 listenkv.Store,由 rootmulti store 用来包装 KVStore,以启用状态监听。
我们将使用 MemoryListener 来配置这个 Store,它会收集状态变化并输出到特定目标。
MultiStore 接口更新
我们将更新CommitMultiStore 接口,以便将 Memorylistener 包装到特定的 KVStore 上。
请注意,MemoryListener 将由具体的 rootmulti 实现在内部附加。
MultiStore 实现更新
我们将调整rootmulti 的 GetKVStore 方法:如果为该 Store 启用了监听,则用 listenkv.Store 包装返回的 KVStore。
AddListeners,在内部管理 KVStore 监听器,并实现 PopStateCache 以提供获取当前状态的方式。
rootmulti 的 CacheMultiStore 和 CacheMultiStoreWithVersion 方法,以在缓存层启用监听。
暴露数据
流式服务
我们将引入一个新的ABCIListener 接口,把它接入 BaseApp,并转发 ABCI 请求与响应,
以便服务能够将状态变化与 ABCI 请求分组关联起来。
BaseApp 注册
我们将向BaseApp 添加一个新方法,以支持注册 StreamingService:
BaseApp 结构体添加两个新字段:
ABCI 事件钩子
我们将修改FinalizeBlock 和 Commit 方法,以便将 ABCI 请求和响应传递给任何已向 BaseApp 注册的流式服务钩子。
Go 插件系统
我们提出一种插件架构,用于加载和运行Streaming 插件及其他类型的实现。我们将引入一个基于 gRPC 的插件系统,用于加载和运行 Cosmos-SDK 插件。该插件系统使用 hashicorp/go-plugin。
每个插件都必须有一个结构体来实现 plugin.Plugin 接口,以及一个 Impl 接口来处理通过 gRPC 传递的消息。
每个插件还必须为 gRPC 服务定义消息协议:
plugin.Plugin 接口有两个方法:Client 和 Server。对于我们的 gRPC 服务,这两个方法分别是 GRPCClient 和 GRPCServer。
Impl 字段保存了用 Go 编写的 baseapp.ABCIListener 接口的具体实现。
注意:这仅用于以 Go 编写的插件实现。
这种插件系统的优势在于,插件作者可以在每个插件内部以适合其用例的方式定义消息协议。
例如,当需要监听状态变更时,可以如下定义 ABCIListener 消息协议(仅用于说明)。
当不需要监听状态变更时,可以从协议中省略 ListenCommit。
Impl(这仅用于以 Go 编写的插件):
(interface{}, error)。
这带来了使用带版本插件的优势,即插件接口和 gRPC 协议可以随着时间演进而变化。
此外,它还支持构建独立插件,以便通过 gRPC 暴露系统的不同部分。
RegisterStreamingPlugin 函数,用于向应用的 BaseApp 注册 NewStreamingPlugin。
Streaming 插件可以是 Any 类型;因此,该函数接收的是接口而不是具体类型。
例如,我们可以有 ABCIListener、WasmListener 或 IBCListener 类型的插件。注意,RegisterStreamingPluing 函数只是辅助函数,并非必需。插件注册很容易从应用侧直接移到 BaseApp 中。
NewStreamingPlugin 和 RegisterStreamingPlugin 函数用于将插件注册到应用的 BaseApp。
例如,在 NewSimApp 中:
配置
插件系统将在应用的 TOML 配置文件中进行配置。Whether to stop the node on message deliver error.
stop-node-on-err = trueListenKVStore
We will create a newStore type listenkv.Store that the rootmulti store will use to wrap a KVStore to enable state listening.
We will configure the Store with a MemoryListener which will collect state changes for output to specific destinations.
MultiStore interface updates
We will update theCommitMultiStore interface to allow us to wrap a Memorylistener to a specific KVStore.
Note that the MemoryListener will be attached internally by the concrete rootmulti implementation.
MultiStore implementation updates
We will adjust therootmulti GetKVStore method to wrap the returned KVStore with a listenkv.Store if listening is turned on for that Store.
AddListeners to manage KVStore listeners internally and implement PopStateCache
for a means of retrieving the current state.
rootmulti CacheMultiStore and CacheMultiStoreWithVersion methods to enable listening in
the cache layer.
Exposing the data
Streaming Service
We will introduce a newABCIListener interface that plugs into the BaseApp and relays ABCI requests and responses
so that the service can group the state changes with the ABCI requests.
BaseApp Registration
We will add a new method to theBaseApp to enable the registration of StreamingServices:
BaseApp struct:
ABCI Event Hooks
We will modify theFinalizeBlock and Commit methods to pass ABCI requests and responses
to any streaming service hooks registered with the BaseApp.
Go Plugin System
We propose a plugin architecture to load and runStreaming plugins and other types of implementations. We will introduce a plugin
system over gRPC that is used to load and run Cosmos-SDK plugins. The plugin system uses hashicorp/go-plugin.
Each plugin must have a struct that implements the plugin.Plugin interface and an Impl interface for processing messages over gRPC.
Each plugin must also have a message protocol defined for the gRPC service:
plugin.Plugin interface has two methods Client and Server. For our GRPC service these are GRPCClient and GRPCServer
The Impl field holds the concrete implementation of our baseapp.ABCIListener interface written in Go.
Note: this is only used for plugin implementations written in Go.
The advantage of having such a plugin system is that within each plugin authors can define the message protocol in a way that fits their use case.
For example, when state change listening is desired, the ABCIListener message protocol can be defined as below (for illustrative purposes only).
When state change listening is not desired than ListenCommit can be omitted from the protocol.
Impl(this is only used for plugins that are written in Go):
(interface{}, error).
This provides the advantage of using versioned plugins where the plugin interface and gRPC protocol change over time.
In addition, it allows for building independent plugin that can expose different parts of the system over gRPC.
RegisterStreamingPlugin function for the App to register NewStreamingPlugins with the App’s BaseApp.
Streaming plugins can be of Any type; therefore, the function takes in an interface vs a concrete type.
For example, we could have plugins of ABCIListener, WasmListener or IBCListener. Note that RegisterStreamingPluing function
is helper function and not a requirement. Plugin registration can easily be moved from the App to the BaseApp directly.
NewStreamingPlugin and RegisterStreamingPlugin functions are used to register a plugin with the App’s BaseApp.
e.g. in NewSimApp:
Configuration
The plugin system will be configured within an App’s TOML configuration files.ABCIListener plugin: streaming.abci.plugin, streaming.abci.keys, streaming.abci.async and streaming.abci.stop-node-on-err.
streaming.abci.plugin is the name of the plugin we want to use for streaming, streaming.abci.keys is a set of store keys for stores it listens to,
streaming.abci.async is bool enabling asynchronous listening and streaming.abci.stop-node-on-err is a bool that stops the node when true and when operating
on synchronized mode streaming.abci.async=false. Note that streaming.abci.stop-node-on-err=true will be ignored if streaming.abci.async=true.
The configuration above support additional streaming plugins by adding the plugin to the [streaming] configuration section
and registering the plugin with RegisterStreamingPlugin helper function.
Note the that each plugin must include streaming.{service}.plugin property as it is a requirement for doing the lookup and registration of the plugin
with the App. All other properties are unique to the individual services.
Encoding and decoding streams
ADR-038 introduces the interfaces and types for streaming state changes out from KVStores, associating this data with their related ABCI requests and responses, and registering a service for consuming this data and streaming it to some destination in a final format. Instead of prescribing a final data format in this ADR, it is left to a specific plugin implementation to define and document this format. We take this approach because flexibility in the final format is necessary to support a wide range of streaming service plugins. For example, the data format for a streaming service that writes the data out to a set of files will differ from the data format that is written to a Kafka topic.Consequences
These changes will provide a means of subscribing to KVStore state changes in real time.Backwards Compatibility
- This ADR changes the
CommitMultiStoreinterface, implementations supporting the previous version of this interface will not support the new one
Positive
- Ability to listen to KVStore state changes in real time and expose these events to external consumers
Negative
- Changes
CommitMultiStoreinterface and its implementations
Neutral
- Introduces additional- but optional- complexity to configuring and running a cosmos application
- If an application developer opts to use these features to expose data, they need to be aware of the ramifications/risks of that data exposure as it pertains to the specifics of their application