rpc4cj 库
介绍
一个高性能、开源和通用的 RPC 框架,基于ProtoBuf(Protocol Buffers) 序列化协议开发,且支持众多开发语言。面向服务端和移动端,基于 HTTP/2 设计,带来诸如双向流、流控、头部压缩、单 TCP 连接上的多复用请求等特。这些特性使得其在移动设备上表现更好,更省电和节省空间占用。
1 一元RPC交互方式客户端与服务端的信息通信功能
前置条件:NA
场景:
1.提供客户端一元请求发送;
2.提供服务端一元请求处理。
约束: NA
性能: 支持版本几何性能持平
可靠性: NA
1.1 客户端一元请求发送
1.1.1 主要接口
public class ClientConn {
/*
* 获取当前 clientconn 目标地址
* 返回值 String - 当前 grpc 的请求目标地址
*/
public func getTarget(): String
/*
* 关闭当前和服务器的连接
*/
public func close()
/*
* 获取当前 grpc 请求状态
* 返回值 ConState - 当前 grpc 请求状态
*/
public func getState(): ConState
/*
* 添加新的地址目标容器 clientconn 的 conn 属性
* 参数 String - 要添加目标地址
*/
public func newAddrConn(addr: String): Unit
/*
* 判断当前状态是否等于 sourceState
* 参数 sourceState - 判断的状态
* 返回值 Bool - 是否相等
*/
public func WaitForStateChange(sourceState: ConState): Bool
/*
* 获取方法的 MethodConfig 信息
* 参数 method - 请求方法名
* 返回值 MethodConfig - 方法配置实例对象
*/
public func GetMethodConfig(method: String): MethodConfig
}
public enum ConState {
|IDLE
|CONNECTING
|READY
|TRANSIENT_FAILURE
|SHUTDOWN
|INVALID_STATE
/*
* 操作符运算 ==
* 参数 that- 用于比较的其他的 ConState
* 参数 Bool - 相等返回 true;否则返回 false
*/
public operator func == (that: ConState): Bool
/*
* 操作符运算 !=
* 参数 that - 用于比较的其他的 ConState
* 参数 Bool - 不相等返回 true;否则返回 false
*/
public operator func != (that: ConState): Bool
}
/*
* 和 grpc 服务器创建连接, 该 api 为 dial 的核心 api
* 参数 target - 服务器地址
* 参数 cfg - 客户端tls配置,默认不使用tls安全协议
* 返回值 ClientConn - 返回一个 ClientConn 实例对象
*/
public func dial(target: String, cfg!: TlsClientConfig = TlsClientConfig("")): ClientConn
/*
* 带额外连接参数的 dial
* 参数 target - 服务器地址
* 参数 opts - 连接参数 DialOption 列表,目前仅支持 withResolvers
* 参数 cfg - 客户端tls配置,默认不使用tls安全协议
* 返回值 ClientConn - 返回一个 ClientConn 实例对象
*/
public func dial(target: String, opts: ArrayList<DialOption>, cfg!: TlsClientConfig = TlsClientConfig("")): ClientConn
/*
* 指定 gRPC 是否因非临时拨号错误而失败
* 参数 bol - 是否启用
* 返回值 DialOption - 返回的 DialOption
*/
public func failOnNonTempDialError(bol: Bool): DialOption
/*
* 根据指定函数创建 FuncDialOption
* 参数 f - (DialOptions) -> Unit 类型函数
* 返回值 FuncDialOption - 返回的 FuncDialOption
*/
public func newFuncDialOption(f: (DialOptions) -> Unit ): FuncDialOption {
return FuncDialOption(f)
}
/*
* 指定要用作 authority 伪标头的值,并在身份验证握手中用作服务器名称
* 参数 str - 设置的值
* 返回值 DialOption - 返回的 DialOption
*/
public func withAuthority(str: String): DialOption
/*
* 设置在连接失败后使用提供的回退参数。
* 参数 b - 设置的 BackoffConfig 参数
* 返回值 DialOption - 返回的 DialOption
*/
public func withBackoffConfig(b: BackoffConfig): DialOption
/*
* 设置用于连接的回退策略失败后重试次数
* 参数 bs - 设置的 Strategy
* 返回值 DialOption - 返回的 DialOption
*/
public func withBackoff(bs: Strategy): DialOption
/*
* 设置读取缓冲区的大小
* 参数 s - 设置的大小
* 返回值 DialOption - 返回的 DialOption
*/
public func withReadBufferSize(s: Int64): DialOption
/*
* 设置流上初始窗口大小的值
* 参数 s - 设置的大小
* 返回值 DialOption - 返回的 DialOption
*/
public func withInitialWindowSize(s: Int32): DialOption
/*
* 设置连接上初始窗口大小的值
* 参数 s - 设置的大小
* 返回值 DialOption - 返回的 DialOption
*/
public func withInitialConnWindowSize(s: Int32): DialOption
/*
* 设置客户端可以接收的最大消息大小
* 参数 s - 设置的大小
* 返回值 DialOption - 返回的 DialOption
*/
public func withMaxMsgSize(s: Int64): DialOption
/*
* 设置用于消息封送处理和取消封送处理的编解码器
* 参数 c - 设置的编解码器
* 返回值 DialOption - 返回的 DialOption
*/
public func withCodec(c: BaseCodec): DialOption
/*
* 设置使拨号的呼叫者阻塞,直到基础连接启动。
* 返回值 DialOption - 返回的 DialOption
*/
public func withBlock(): DialOption
/*
* 确定在网络上执行写入操作之前可以批量处理的数据量
* 参数 s - 设置的数据量大小
* 返回值 DialOption - 返回的 DialOption
*/
public func withWriteBufferSize(s: Int64): DialOption
/*
* 指定所有 RPC 的用户代理字符串
* 参数 s - 指定的字符串
* 返回值 DialOption - 返回的 DialOption
*/
public func withUserAgent(s: String): DialOption
/*
* 配置客户端的超时等待时间
* 参数 d - 设置等待时间长度
* 返回值 DialOption - 返回的 DialOption
*/
public func withTimeout(d: Duration): DialOption
/*
* 使客户端连接返回一个字符串,其中包含最后一个连接错误
* 返回值 DialOption - 返回的 DialOption
*/
public func withReturnConnectionError(): DialOption
/*
* 禁止使用此客户端通信的代理
* 返回值 DialOption - 返回的 DialOption
*/
public func withNoProxy(): DialOption
/*
* 指定客户端准备的头列表的最大大小
* 参数 s - 指定的大小
* 返回值 DialOption - 返回的 DialOption
*/
public func withMaxHeaderListSize(s: UInt32): DialOption
/*
* 禁用此客户端通讯子的传输安全性
* 返回值 DialOption - 返回的 DialOption
*/
public func withInsecure(): DialOption
/*
* 使 gRPC 忽略解析程序提供的任何服务配置
* 返回值 DialOption - 返回的 DialOption
*/
public func withDisableServiceConfig(): DialOption
/*
* 禁用重试
* 返回值 DialOption - 返回的 DialOption
*/
public func withDisableRetry(): DialOption
/*
* 禁用此客户端康纳的所有子计算机的 LB 通道运行状况检查
* 返回值 DialOption - 返回的 DialOption
*/
public func withDisableHealthCheck(): DialOption
/*
* 配置默认服务配置
* 参数 s - 配置的字符串
* 返回值 DialOption - 返回的 DialOption
*/
public func withDefaultServiceConfig(s: String): DialOption
/*
* 配置解析器
* 参数 rs - 配置的解析器
* 返回值 DialOption - 返回的 DialOption
*/
public func withResolvers(rs: resolver.Builder): DialOption
/*
* 配置一元请求拦截器
* 参数 unary - 配置的一元拦截器
* 返回值 DialOption - 返回的 DialOption
*/
public func WithUnaryInterceptor(unary: UnaryInvoker): DialOption
/*
* 配置流式请求拦截器
* 参数 unary - 配置的流式拦截器
* 返回值 DialOption - 返回的 DialOption
*/
public func withStreamInterceptor(unary: StreamClientInterceptor): DialOption
/*
* 返回一个 CallOption,设置客户端可以接收的最大消息大小
* 参数 m - 设置的大小
* 返回值 CallOption - 返回的 CallOption
*/
public func maxCallRecvMsgSize(bytes: Int64): CallOption
/*
* 返回一个 CallOption,该选项将设置用于调用的所有请求和响应消息的编解码器
* 参数 m - 设置的 BaseCodec
* 返回值 CallOption - 返回的 CallOption
*/
public func callCustomCodec(codec: BaseCodec): CallOption
/*
* 为一个请求设置 content-subtype
* 参数 m - 设置的 contentSubtype 字符串
* 返回值 CallOption - 返回的 CallOption
*/
public func CallContentSubtype(contentSubtype: String): CallOption
/*
* 返回一个 CallOption,该选项将设置用于调用的所有请求和响应消息的编解码器
* 参数 m - 设置的 BaseCodec
* 返回值 CallOption - 返回的 CallOption
*/
public func forceCodec(codec: BaseCodec): CallOption
/*
* 返回一个 CallOption,设置客户端可以发送的最大消息大小
* 参数 m - 设置的消息大小
* 返回值 CallOption - 返回的 CallOption
*/
public func maxCallSendMsgSize(bytes: Int64): CallOption
/*
* 限制用于缓冲此 RPC 的请求进行重试的内存量。
* 参数 m - 设置的重试的内存量大小
* 返回值 CallOption - 返回的 CallOption
*/
public func maxRetryRPCBufferSize(bytes: Int64): CallOption
/*
* 设置发送请求时使用的压缩器
* 参数 m - 设置使用的压缩器
* 返回值 CallOption - 返回的 CallOption
*/
public func useCompressor(name: String): CallOption
/*
* 配置在断开的连接或无法访问的服务器上尝试 RPC 时要执行的操作
* 参数 m - 设置是否启用
* 返回值 CallOption - 返回的 CallOption
*/
public func waitForReady(waitForReady: Bool): CallOption
/*
* 配置 PerRPCCredentials 指示用于呼叫的每个 RPC 凭据
* 参数 m - 设置的 PerRPCCredentials
* 返回值 CallOption - 返回的 CallOption
*/
public func perRPCCredentials(creds: PerRPCCredentials): CallOption
/**
* 设置发起 grpc 请求后等待服务器响应的时长,未设置时默认等待2小时。
* 参数 ns - 设置等待的时间,单位纳秒
*/
public func setGrpcTimeout(ns: Int64): Unit {
}
/**
* 设置 grpc 响应的缓存有效期,未设置时默认缓存有效期10000000ns。
* 参数 ns - 设置缓存有效期,单位纳秒
*/
public func setRefreshTime(ns: Int64): Unit
/**
* 获取当前 grpc 响应的缓存有效期。
* 返回值 Int64 - 当前 grpc 响应的缓存有效期,单位纳秒
*/
public func getRefreshTime(): Int64
1.2 服务端一元请求处理
1.2.1 主要接口
1.2.1.1 class Server
服务器是用于服务RPC请求的gRPC服务器。
public class Server {
/**
* 根据可选择参数,创建 Server 实例。
* 参数 opts - Server的服务器选项设置类
* 参数 lis - HashMap 数据结构, key 为 GRPCServerSocket, value 为 Bool 类型
* 参数 conns - 包含所有活动服务器传输。它是一个映射,key 为监听地址,值是属于该侦听器的活动传输集。
* 参数 server - Bool 数据类型,是否已经启动服务
* 参数 drain - Bool 数据类型,是否已经排空
* 参数 services - HashMap 数据结构, key 为 service 名称, value 为有关服务的信息类 ServiceInfoInside 数据类型
* 参数 czData - 用于存储ClientConn、addrConn和Server的channelz相关数据
*/
public init(opts!: ServerOptions = ServerOptions(),
lis!: HashMap<GRPCServerSocket, Bool> = HashMap<GRPCServerSocket, Bool>(),
conns!: HashMap<String, HashMap<ServerTransport, Bool>> = HashMap<String, HashMap<ServerTransport, Bool>>(),
server!: Bool = false,
drain!: Bool = false,
services!: HashMap<String, ServiceInfoInside> = HashMap<String, ServiceInfoInside>(),
czData!: ChannelzDataInside = ChannelzDataInside())
/**
* 创建了一个 gRPC 服务器,该服务器没有注册服务,也没有开始接受请求。
* 参数 opts - Array 数据结构, 服务器选项设置诸如凭据、编解码器和保留活动参数等选项
* 返回值 Server - Server 实例
*/
public static func newServer(opt: Array<ServerOption>): Server
/**
* 向gRPC服务器注册服务及其实现。这必须在调用服务之前调用。
* 参数 desc - RPC 服务的规范
* 参数 ss - MessageLite 数据结构,可以是实际应用中自行实现的protobuf实现类
* 注意:参数 ss 在当前版本支持的参数类型为MessageLite,后续版本更迭会修改参数类型为Any
*/
public func registerService(desc: ServiceDesc, ss: MessageLite): Unit
/**
* 返回从服务名称到服务信息的映射。服务名称包括包名称,形式为<package>、<Service>。
* 返回值 HashMap<String, ServiceInfo> - HashMap 数据结构, key 为 service 名称, value 为 ServiceInfo 数据类型
*/
public func getServiceInfo(): HashMap<String, ServiceInfo>
/**
* 接受监听器lis上的传入连接,为每个连接创建一个新的ServerTransport和service 数据处理线程。
* 在数据处理线程中,读取gRPC请求,然后调用注册的处理程序来回复它们。
* 当lis接受失败并出现致命错误时,Serve返回。此方法返回时,lis将关闭。
* 除非调用Stop或GracefulStop,否则Serve将返回非零错误。
* 参数 ss - GRPCServerSocket 数据结构, 成员变量 ss 中封装 SocketServer/TlsSocketServer
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public func serve(ss: GRPCServerSocket): GrpcError
/**
* ServeHTTP通过响应gRPC请求r,通过在gRPC服务器s中查找请求的gRPC方法,实现了Go标准库的http处理程序接口。
* 提供的HTTP请求必须通过HTTP/2连接到达。当使用Go标准库的服务器时,实际上这意味着请求也必须通过TLS到达。
* 要在gRPC和现有http处理程序之间共享一个端口(如https的443),请使用根http处理程序,如:
* if r.ProtoMajor == 2 && strings.HasPrefix(
* r.Header.Get("Content-Type"), "application/grpc") {
* grpcServer.ServeHTTP(w, r)
* } else {
* yourMux.ServeHTTP(w, r)
* }
* 请注意,ServeHTTP使用Go的HTTP/2服务器实现,它与grpc Go的HTTP/2服务器完全分离。
* 两条路径之间的性能和特性可能有所不同。ServeHTTP不支持通过GRPCGo的HTTP/2服务器提供的一些gRPC功能。
*
* 注意:此API是实验性的,可能会在以后的版本中更改或删除。
*
* 参数 w - ResponseWriteStream 接口,用于处理响应信息。
* 参数 r - Request 数据结构
* 备注:参数 r 必须同时符合以下3个条件:
* 1. Request.protoMajor = 2;
* 2. 必须为 POST 请求;
* 3. Request.header 必须包含 key = "Content-Type", 且 value 值的第一个为 "application/grpc" 的键值对。
* 备注:若参数 r 未同时符合以上3个条件,代码逻辑跳过本函数,对外不感知:
*/
public func serveHTTP(w: ResponseWriteStream, r: Request): Unit
/**
* 停止停止gRPC服务器。它会立即关闭所有打开的连接和侦听器。
* 它将取消服务器端的所有活动RPC,客户端上相应的挂起RPC将收到连接错误通知。
*/
public func stop(): Unit
/**
* 优雅地停止gRPC服务器。它将停止服务器接受新的连接、RPC和块,直到所有挂起的RPC完成。
*/
public func gracefulStop(): Unit
}
1.2.1.2 interface ServerOption
服务器选项设置诸如凭据、编解码器和保留活动参数等选项。
public interface ServerOption {
/**
* 传入参数so,执行apply函数。
* 参数 so - Server的服务器选项设置类
*/
func apply(so: ServerOptions): Unit
}
1.2.1.3 class MethodDesc
表示RPC服务的方法规范。
public class MethodDesc {
/**
* 根据可选择参数,创建 MethodDesc 实例。
* 参数 methodName - String 数据结构
* 参数 handler - 函数 数据结构
*/
public init(methodName!: String = "",
handler!: (MessageLite, (MessageLite) -> GrpcError, ?UnaryServerInterceptor) -> (MessageLite, GrpcError))
}
1.2.1.4 class MethodInfo
包含有关RPC的信息,包括其方法名称和类型。
public class MethodInfo {
/**
* 根据可选择参数,创建 MethodInfo 实例。
* 参数 name - String 数据结构,该名称只是方法名,不包含服务名或包名。
* 参数 isClientStream - Bool 数据结构,指示RPC是否是客户端流式RPC。
* 参数 isServerStream - Bool 数据结构,指示RPC是否是服务器流式RPC。
*/
public init(name!: String = "", isClientStream!: Bool = false, isServerStream!: Bool = false)
}
1.2.1.5 class ServiceInfo
包含服务的一元RPC方法信息、流式RPC方法信息和元数据。
public class ServiceInfo {
/**
* 根据可选择参数,创建 ServiceInfo 实例。
* 参数 methods - Array<MethodInfo> 数据结构
* 参数 metadata - Any 数据结构,元数据是注册服务时在ServiceDesc中指定的元数据。
*/
public init(methods!: Array<MethodInfo> = Array<MethodInfo>(), metadata!: Any)
}
1.2.1.6 class UnaryHandler
定义由UnaryServerInterceptor调用的处理程序,以完成一元RPC的正常执行。
如果一个UnaryHandler返回一个错误,它应该由状态包生成,或者是上下文错误之一。
否则,gRPC将使用Unknown作为状态代码,使用Error()作为RPC的状态消息。
public class UnaryHandler {
/**
* 根据参数函数f,创建 StreamServerInterceptor 实例。
* 参数 f - 函数类型
*/
public init(f: (Any) -> (Any, GrpcError))
}
1.2.1.7 class UnaryServerInterceptor
提供了一个钩子来拦截服务器上一元RPC的执行。info包含拦截器可以操作的这个RPC的所有信息。
handler是服务方法实现的包装器。拦截器负责调用处理程序来完成RPC。
public class UnaryServerInterceptor {
/**
* 根据参数函数f,创建 StreamServerInterceptor 实例。
* 参数 f - 函数类型
*/
public init(f: (Any, UnaryServerInfo, UnaryHandler) -> (MessageLite, GrpcError))
}
1.2.1.8 class UnaryServerInfo
包含关于服务器端一元RPC的各种信息。拦截器可以改变所有每rpc信息。
public class UnaryServerInfo {
/**
* 根据可选择参数,创建 UnaryServerInfo 实例。
* 参数 server - Any 数据结构,服务器是用户提供的服务实现。这是只读的。
* 参数 fullMethod - String 数据结构,是完整的RPC方法字符串,即/package.service/method。
*/
public init(server!: Any, fullMethod!: String = "")
}
1.2.1.9 class ServerConfig
包含建立服务器传输的所有配置。
public class ServerConfig {
/**
* 根据可选择参数,创建 ServerConfig 实例。
* 参数 maxStreams - UInt32 数据结构
* 参数 connectionTimeout - Duration 数据结构
* 参数 inTapHandle - ServerInHandle 数据结构
* 参数 statsHandlers - ArrayList<GrpcHandler> 数据结构
* 参数 keepaliveParams - ServerParameters 数据结构
* 参数 keepalivePolicy - EnforcementPolicy 数据结构
* 参数 initialWindowSize - Int32 数据结构
* 参数 initialConnWindowSize - Int32 数据结构
* 参数 writeBufferSize - Int64 数据结构
* 参数 readBufferSize - Int64 数据结构
* 参数 maxHeaderListSize - UInt32 数据结构
* 参数 headerTableSize - UInt32 数据结构
*/
public init(
maxStreams !: UInt32 = 0,
connectionTimeout !: Duration = Duration.second(),
inTapHandle !: ServerInHandle = ServerInHandle(),
statsHandlers !: ArrayList<GrpcHandler> = ArrayList<GrpcHandler>(),
keepaliveParams !: ServerParameters = ServerParameters(),
keepalivePolicy !: EnforcementPolicy = EnforcementPolicy(),
initialWindowSize !: Int32 = 0,
initialConnWindowSize !: Int32 = 0,
writeBufferSize !: Int64 = 0,
readBufferSize !: Int64 = 0,
maxHeaderListSize !: UInt32 = 0,
headerTableSize !: UInt32 = 0)
}
1.2.1.10 class Stream
表示传输层中的RPC。
public open class Stream <: InputStream {
/**
* 根据可选择参数,创建 Stream 实例。
* 参数 id - UInt32 数据结构,流 id
* 参数 method - String 数据结构,mothod 名称
* 参数 recvCompress - String 数据结构,接收压缩名称
* 参数 st - Option<ServerTransport> 数据结构,ServerTransport 是所有 gRPC 服务器端传输实现的通用父类
* 参数 fc - InFlow 数据结构,流入处理入站流量控制
* 参数 contentSubtype - String 数据结构,是请求的内容子类型
* 参数 isNone - Bool 数据结构,创建的实例是否是空结构
*/
public init(id!: UInt32 = 0,
method! : String = "",
recvCompress! : String = "",
st!: ?ServerTransport = None,
fc!: InFlow = NULL_INFLOW,
contentSubtype!: String = "",
isNone!: Bool = false)
/**
* 从输入流中读取数据放到 data 中
* 参数 data - 读取数据存放的缓冲区,若 buffer 为空则抛出异常
* 返回值 Int64 - 读取成功,返回读取字节数
* 若流被关闭或者没有数据可读,则返回 0
* 读取失败,则抛出异常
*/
public func read(data: Array<Byte>): Int64
/**
* 返回应用于入站消息的压缩算法。如果未应用压缩,则为空字符串。
* 返回值 String - 返回应用于入站消息的压缩算法名称
*/
public func getRecvCompress(): String
/**
* 拖车返回缓存的拖车metedata。请注意,如果在完成整个流之后不调用它,它可能只返回空的MD.客户端。
* 只有在流结束后,即读或写返回io.EOF,它才能安全地读取。
* 返回值 HashMap<String, Array<String>> - 返回拖车返回缓存的拖车metedata
*/
public func getTrailer(): HashMap<String, Array<String>>
/**
* 返回请求的内容子类型。
* 返回值 String - 返回请求的内容子类型名称
*/
public func getContentSubtype(): String
/**
* 返回流的方法。
* 返回值 String - 返回流的方法名称
*/
public func getMethod(): String
/**
* 状态返回从服务器接收的状态。只有在流结束后,也就是在完成关闭后,才能安全读取状态。
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public func getStatus(): GrpcError
/**
* 设置头元数据。这可以被称为多次。仅限服务器端。这不应与其他数据写入并行调用。
* 参数 md - HashMap<String, Array<String>>,头部的键值对数据
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public open func setHeader(md: HashMap<String, Array<String>>): GrpcError
/**
* 发送给定的头元数据。给定元数据与之前调用SetHeader设置的任何元数据组合,然后写入传输流。
* 参数 md - HashMap<String, Array<String>>,头部的键值对数据
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public open func sendHeader(md: HashMap<String, Array<String>>): GrpcError
/**
* 设置服务器将随RPC状态一起发送的尾部元数据。这可以调用多次。仅服务器端。
* 这不应与其他数据写入并行调用。
* 参数 md - HashMap<String, Array<String>>,头部的键值对数据
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public open func setTrailer(md: HashMap<String, Array<String>>): GrpcError
}
1.2.1.11 class ServerTransport
是所有gRPC服务器端传输实现的通用接口。
public abstract class ServerTransport <: Hashable & Equatable<ServerTransport> {
/**
* 使用给定的处理程序接收传入流。
* 参数 handle - 函数类型
* 参数 traceCtx - 函数类型
*/
public func handleStreams(handle: (Stream) -> Unit, traceCtx: (String) -> Unit): Unit
/**
* 发送给定流的头元数据。不能对所有流调用WriteHeader。
* 参数 s - Stream 数据结构
* 参数 md - HashMap 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public func writeHeader(s: Stream, md: HashMap<String, Array<String>>): GrpcError
/**
* 发送给定流的数据。不能对所有流调用Write。
* 参数 s - Stream 数据结构
* 参数 hdr - Array<UInt8> 数据结构
* 参数 data - Array<UInt8> 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public func write(s: Stream, hdr: Array<UInt8>, data: Array<UInt8>): GrpcError
/**
* 将流的状态发送给客户端。WriteStatus是对流进行的最后一个调用,并且总是发生。
* 参数 s - Stream 数据结构
* 参数 st - GrpcError 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public func writeStatus(s: Stream, st: GrpcError): GrpcError
/**
* 关闭会撕裂传输线。一旦调用,就不应再访问传输。所有挂起的流及其处理程序将异步终止。
*/
public func close(): Unit
/**
* 返回远程网络地址。
* 返回值 SocketAddress - SocketAddress 数据结构,返回远程网络地址。
*/
public func getRemoteAddr(): SocketAddress
/**
* 通知客户端此ServerTransport停止接受新的RPC。
*/
public func drain(): Unit
/**
* 增加通过此传输发送的消息数。
*/
public func incrMsgSent(): Unit
/**
* 增加通过此传输接收的消息数。
*/
public func incrMsgRecv(): Unit
}
1.2.1.12 interface GrpcHandler
定义相关统计处理的接口(例如,RPC、连接)。
public interface GrpcHandler {
/**
* 传入参数info,执行函数。
* 参数 info - RPCTagInfo 数据结构,定义 RPC 上下文标记器所需的相关信息
*/
func tagRPC(info: RPCTagInfo): Unit
/**
* 处理RPC统计数据。
* 参数 stats - RPCStats 数据结构,包含关于RPC的统计信息。
*/
func handleRPC(stats: RPCStats): Unit
/**
* 可以将一些信息附加到给定的上下文中。返回的上下文将用于统计处理。对于conn stats处理,HandleConn中用于此连接的上下文将从返回的上下文派生。
* 对于RPC状态处理,
* -在服务器端,HandleRPC中用于此连接上所有RPC的上下文将从返回的上下文派生。
* -在客户端,上下文不是从返回的上下文派生的。
* 参数 ctInfo - ConnTagInfo 数据结构
*/
func tagConn(ctInfo: ConnTagInfo): Unit
/**
* 处理Conn统计数据。
* 参数 cs - ConnStats 数据结构
*/
func handleConn(cs: ConnStats): Unit
}
1.2.1.13 class ConnTagInfo
定义连接上下文标记器所需的相关信息。
public class ConnTagInfo {
/**
* 根据可选择参数,创建 ConnTagInfo 实例。
* 参数 remoteAddr - SocketAddress 数据类型,是相应连接的远程地址
* 参数 localAddr - SocketAddress 数据类型,是相应连接的本地地址。
*/
public init(remoteAddr!: SocketAddress = DEFAULT_SADDRESS, localAddr!: SocketAddress = DEFAULT_SADDRESS)
}
1.2.1.14 class RPCTagInfo
定义RPC上下文标记器所需的相关信息。
public class RPCTagInfo {
/**
* 根据可选择参数,创建 RPCTagInfo 实例。
* 参数 fullMethodName - String 数据类型,package.service/method格式的RPC方法。
* 参数 failFast - Bool 数据类型,此RPC是否为FailFat。此字段仅在客户端有效,在服务器端始终为false。
*/
public init (fullMethodName!: String = "", failFast!: Bool = false) {
this.fullMethodName = fullMethodName
this.failFast = failFast
}
}
1.2.1.15 interface RPCStats
包含关于RPC的统计信息。
public interface RPCStats {
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
func isClient(): Bool
}
1.2.1.16 class Begin
接口RPCStats的实现类,包含RPC尝试开始时的统计信息。FailFast仅在从客户端开始时有效。
public class Begin <: RPCStats {
/**
* 根据可选择参数,创建 Begin 实例。
* 参数 client - 如果来自客户端,则为 true。
* 参数 beginTime - Time 数据结构, 是RPC尝试开始的时间
* 参数 failFast - Bool 数据类型,指示此RPC是否为FailFast
* 参数 isClientStream - Bool 数据类型,指示RPC是否是客户端流式RPC
* 参数 isServerStream - Bool 数据类型,指示RPC是否为服务器流式RPC
* 参数 isTransparentRetryAttempt - Bool 数据结构, 指示此尝试是否是由于透明地重试以前的尝试而启动的
*/
public init(client!: Bool = false,
beginTime!: Time = Time.now(),
failFast!: Bool = false,
isClientStream!: Bool = false,
isServerStream!: Bool = false,
isTransparentRetryAttempt!: Bool = false)
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.17 class InPayload
接口RPCStats的实现类,包含传入有效负载的信息。
public class InPayload <: RPCStats {
/**
* 根据可选择参数,创建 InPayload 实例。
* 参数 client - 如果来自客户端,则客户端为真
* 参数 payload - Any 数据结构,原始类型的有效载荷
* 参数 data - Array<UInt8> 数据类型,序列化消息负载
* 参数 length - Int64 数据类型,未压缩数据的长度
* 参数 wireLength - Int64 数据类型,有线数据的长度(压缩、签名、加密)
* 参数 recvTime - Time 数据结构,收到有效载荷的时间
*/
public init(client!: Bool = false,
payload!: Any,
data!: Array<UInt8> = EMPTY_ARRAY_UINT8,
length!: Int64 = 0,
wireLength!: Int64 = 0,
recvTime!: Time = Time.now())
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.18 class InHeader
接口RPCStats的实现类,包含收到报头时的统计信息。
public class InHeader <: RPCStats {
/**
* 根据可选择参数,创建 InHeader 实例。
* 参数 client - 如果来自客户端,则客户端为真
* 参数 wireLength - Int64 数据类型,头部的有线长度
* 参数 compression - String 数据类型,用于RPC的压缩算法
* 参数 header - HashMap<String, Array<String>> 数据类型,包含接收到的标头元数据
* 参数 fullMethod - String 数据结构,是完整的RPC方法字符串,即/package.service/method
* 参数 remoteAddr - SocketAddress 数据类型,相应连接的远程地址
* 参数 localAddr - SocketAddress 数据结构,相应连接的本地地址
*/
public init(client!: Bool = false, wireLength!: Int64 = 0, compression!: String = "",
header!: HashMap<String, Array<String>> = HashMap<String, Array<String>>(),
fullMethod!: String = "", remoteAddr!: SocketAddress = EMPTY_ADDR,
localAddr!: SocketAddress = EMPTY_ADDR)
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.19 class InTrailer
接口RPCStats的实现类,包含接收到拖车时的统计信息。
public class InTrailer <: RPCStats {
/**
* 根据可选择参数,创建 InTrailer 实例。
* 参数 client - 如果来自客户端,则客户端为真
* 参数 wireLength - Int64 数据类型,trailer 的长度。
* 参数 trailer - HashMap<String, Array<String>> 数据结构,拖车包含从服务器接收的拖车元数据。仅当此InTrailer来自客户端时,此字段才有效。
*/
public init(client!: Bool = false, wireLength!: Int64 = 0,
trailer!: HashMap<String, Array<String>> = HashMap<String, Array<String>>())
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.20 class OutPayload
接口RPCStats的实现类,包含传出有效负载的信息。
public class OutPayload <: RPCStats {
/**
* 根据可选择参数,创建 OutPayload 实例。
* 参数 client - 如果来自客户端,则客户端为真
* 参数 payload - Any 数据结构,原始类型的有效载荷
* 参数 data - Array<UInt8> 数据类型,序列化消息负载
* 参数 length - Int64 数据类型,未压缩数据的长度
* 参数 wireLength - Int64 数据类型,有线数据的长度(压缩、签名、加密)
* 参数 sentTime - Time 数据结构,发送有效载荷的时间
*/
public init(client!: Bool = false,
payload!: Any,
data!: Array<UInt8> = EMPTY_ARRAY_UINT8,
length!: Int64 = 0,
wireLength!: Int64 = 0,
sentTime!: Time = Time.now()
)
/**
* 根据可选择参数,创建 OutPayload 实例。
* 参数 client - 如果来自客户端,则客户端为真
* 参数 msg - Any 数据结构,原始类型的有效载荷
* 参数 data - Array<UInt8> 数据类型,序列化消息负载;data.size 为 length
* 参数 payload - Array<UInt8> 数据类型,payload + HEADER_LEN(5) = wireLength
* 参数 t - Time 数据结构,发送有效载荷的时间
*/
public static func getOutPayload(client: Bool, msg: Any, data: Array<UInt8>, payload: Array<UInt8>, t: Time): OutPayload
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.21 class OutHeader
接口RPCStats的实现类,发送标头时包含统计信息。
public class OutHeader <: RPCStats {
/**
* 根据可选择参数,创建 OutHeader 实例。
* 参数 client - 如果来自客户端,则客户端为真
* 参数 compression - String 数据类型,用于RPC的压缩算法
* 参数 header - HashMap<String, Array<String>> 数据类型,包含发送的标头元数据
* 参数 fullMethod - String 数据结构,是完整的RPC方法字符串,即/package.service/method
* 参数 remoteAddr - SocketAddress 数据类型,相应连接的远程地址
* 参数 localAddr - SocketAddress 数据结构,相应连接的本地地址
*/
public init(client!: Bool = false, compression!: String = "",
header!: HashMap<String, Array<String>> = HashMap<String, Array<String>>(),
fullMethod!: String = "", remoteAddr!: SocketAddress = EMPTY_ADDR,
localAddr!: SocketAddress = EMPTY_ADDR)
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.22 class OutTrailer
接口RPCStats的实现类,包含发送预告片时的统计信息。
public class OutTrailer <: RPCStats {
/**
* 根据可选择参数,创建 OutTrailer 实例。
* 参数 client - 如果来自客户端,则客户端为真
* 参数 wireLength - Int64 数据类型,trailer 的长度。
* 参数 trailer - HashMap<String, Array<String>> 数据结构,拖车包含发送给客户端的拖车元数据。此字段仅在该输出尾文件来自服务器端时有效。
*/
public init(client!: Bool = false, wireLength!: Int64 = 0,
trailer!: HashMap<String, Array<String>> = HashMap<String, Array<String>>())
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.23 class End
接口RPCStats的实现类,包含RPC结束时的统计信息。
public class End <: RPCStats {
/**
* 根据可选择参数,创建 End 实例。
* 参数 client - 如果来自客户端,则为 true。
* 参数 beginTime - Time 数据结构,RPC 尝试开始的时间
* 参数 endTime - Time 数据结构,RPC 结束的时间。
* 参数 trailer - HashMap<String, Array<String>> 数据类型
* 参数 error - GrpcError 数据类型
*/
public init(client!: Bool = false,
beginTime!: Time = Time.now(),
endTime!: Time = Time.now(),
trailer!: HashMap<String, Array<String>> = HashMap<String, Array<String>>(),
error!: GrpcError = NULL_ERR)
/**
* 如果此RPCStats来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果此RPCStats来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.24 interface ConnStats
包含关于连接的统计信息。
public interface ConnStats {
/**
* 如果来自客户端,则返回true。
* 返回值 Bool - Bool 数据结构,如果来自客户端,则返回true。
*/
func isClient(): Bool
}
1.2.1.25 class ConnBegin
接口ConnStats的实现类,包含连接建立时的状态。
public class ConnBegin <: ConnStats {
/**
* 根据可选择参数,创建 ConnBegin 实例。
* 参数 client - 如果来自客户端,则客户端为真
*/
public init(client!: Bool = false)
/**
* 是否来自客户端.
* 返回值 Bool - Bool 数据结构,如果来自客户端,则返回true
*/
public func isClient(): Bool
}
1.2.1.26 class ConnEnd
接口ConnStats的实现类,包含连接结束时的状态。
public class ConnEnd <: ConnStats {
/**
* 根据可选择参数,创建 ConnEnd 实例。
* 参数 client - 如果来自客户端,则客户端为真
*/
public init(client!: Bool = false)
/**
* 是否来自客户端.
* 返回值 Bool - Bool 数据结构,如果来自客户端,则返回 true
*/
public func isClient(): Bool
}
1.2.1.27 class ServerParameters
用于在服务器端设置keepalive和max-age参数。
public class ServerParameters {
/**
* 根据可选择参数,创建 ServerParameters 实例。
* 参数 maxConnectionIdle - Duration 数据结构,空闲持续时间是从最近一次未完成的RPC数变为零或连接建立起定义的。当前默认值为无穷大。
* 参数 maxConnectionAge - Duration 数据结构,连接在通过发送GoAway关闭之前可能存在的最长时间。
* 参数 maxConnectionAgeGrace - Duration 数据结构,是MaxConnectionAge之后的一个加法时段,在此时段之后,连接将被强制关闭。
* 参数 time - Duration 数据结构,经过一段时间后,如果服务器没有看到任何活动,它将ping客户端以查看传输是否仍然有效。如果设置为低于1s,则将使用最小值1s。当前默认值为2小时。
* 参数 timeout - Duration 数据结构,在ping以进行keepalive检查之后,服务器将等待一段时间的超时,如果在此之后没有看到任何活动,则连接将关闭。当前默认值为20秒。
*/
public init(maxConnectionIdle!: Duration = EMPTY_TIME,
maxConnectionAge!: Duration = EMPTY_TIME,
maxConnectionAgeGrace!: Duration = EMPTY_TIME,
time!: Duration = Duration.hour(2),
timeout!: Duration = Duration.second(20))
}
1.2.1.28 class EnforcementPolicy
用于在服务器端设置keepalive强制策略。服务器将关闭与违反此策略的客户端的连接。
public class EnforcementPolicy {
/**
* 根据可选择参数,创建 EnforcementPolicy 实例。
* 参数 minTime - Duration 数据结构,是客户端在发送保活 ping 之前应等待的最短时间。当前默认值为5分钟
* 参数 permitWithoutStream - Bool 数据类型,如果为true,则即使在没有活动流(RPC)的情况下,服务器也允许保持活动ping。
* 如果为false,并且客户端在没有活动流时发送ping,服务器将发送GOAWAY并关闭连接。
* 默认为false。
*/
public init(minTime!: Duration = Duration.minute(5), permitWithoutStream!: Bool = false)
}
1.2.1.29 class Peer
对等体包含RPC对等体的信息,如地址和身份验证信息。
public class Peer {
/**
* 根据可选择参数,创建 Peer 实例。
* 参数 addr - SocketAddress 数据结构,对等地址
* 参数 authInfo - Option<AuthInfo> 数据类型,传输的身份验证信息。如果没有使用传输安全,则为None。
*/
public init(addr!: SocketAddress = EMPTY_ADDR, authInfo!: Option<AuthInfo> = None)
}
1.3 示例
from std import socket.*
from rpc4cj import grpc.*
from rpc4cj import transport.*
from rpc4cj import exceptions.*
from rpc4cj import util.*
from std import time.*
from std import sync.*
from std import collection.*
from protobuf import protobuf.*
let strServer: String = "test pass!"
class ServerTest <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
let hy: HelloReply = HelloReply()
hy.message = strServer
return (hy, GrpcError())
}
public func mustEmbedUnimplementedGreeterServer() {}
}
var port: UInt16 = 0
main() {
let ss: SocketServer = SocketServer(TCP, "127.0.0.1", port)
port = ss.port
spawn{ =>
let lis: GRPCServerSocket = GRPCServerSocket(GrpcSocketS(ss))
let server: Server = Server()
registerGreeterServer(server, ServerTest())
server.serve(lis)
}
sleep(Duration.second)
var req = HelloRequest()
var resp = HelloReply()
req.name = "World"
var con = dial("127.0.0.1:${port}")//linux
con.invoke("/helloworld.Greeter/SayHello", req, resp)
if (!resp.message.contains(strServer)) {
println(-1)
return -1
}
println("test pass")
return 0
}
public class HelloRequest <: MessageLite & TypedMessage<HelloRequest> {
private var p_name = Pb3String()
public init() {
m_innerInit()
}
public mut prop name: String { get() { p_name.get() } set(i_vtmp) { p_name.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_name.isEmpty()) { i_tmp += 1 + binSize(p_name.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_name.isEmpty()
}
protected func m_clear(): Unit {
p_name.clear()
}
public func copyFrom(src: HelloRequest): HelloRequest {
p_name = src.p_name
m_innerCopyFrom(src)
return this
}
public func clone(): HelloRequest {
return HelloRequest().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_name.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_name.isEmpty()) { out.append(10); out.packBin(p_name.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_name.isEmpty()) { i_tmp.append("name", p_name.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloRequest() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
}
public class HelloReply <: MessageLite & TypedMessage<HelloReply> {
private var p_message = Pb3String()
public init() {
m_innerInit()
}
public mut prop message: String { get() { p_message.get() } set(i_vtmp) { p_message.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_message.isEmpty()) { i_tmp += 1 + binSize(p_message.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_message.isEmpty()
}
protected func m_clear(): Unit {
p_message.clear()
}
public func copyFrom(src: HelloReply): HelloReply {
p_message = src.p_message
m_innerCopyFrom(src)
return this
}
public func clone(): HelloReply {
return HelloReply().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_message.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_message.isEmpty()) { out.append(10); out.packBin(p_message.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_message.isEmpty()) { i_tmp.append("message", p_message.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloReply() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
}
//server 端
public abstract class GreeterServer <: MessageLite {
public func SayHello(hr: HelloRequest): (HelloReply, GrpcError)
public func mustEmbedUnimplementedGreeterServer(): Unit
protected func m_size(): Int64 { return 0 }
protected func m_isEmpty(): Bool { return false }
protected func m_clear(): Unit {}
public func unpacki(_: BytesReader): Unit {}
public func packi(_: BytesWriter): Unit {}
public func toString(): String { return "" }
}
class UnimplementedGreeterServer <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
return (HelloReply(), GrpcError(Unimplemented, "method SayHello not implemented"))
}
public func mustEmbedUnimplementedGreeterServer() {}
}
interface UnsafeGreeterServer {
func mustEmbedUnimplementedGreeterServer(): Unit
}
func Greeter_SayHello_Handler(srv: MessageLite, dec: (MessageLite) -> GrpcError, interceptor: ?UnaryServerInterceptor): (MessageLite, GrpcError) {
let inhr: HelloRequest = HelloRequest()
var err: GrpcError = dec(inhr)
if (!err.isNull()) {
return (inhr, err)
}
match(interceptor){
case Some(v) =>
let info: UnaryServerInfo = UnaryServerInfo(
server: srv,
fullMethod: "/helloworld.Greeter/SayHello"
)
println("--------------------------000")
println("inhr.name = ${inhr.name}")
println("--------------------------000")
inhr.name = "dec(inhr)解析到的HelloRequest是空的???-----000"
func handlercs(req: Any): (Any, GrpcError) {
return (srv as (GreeterServer)).getOrThrow().SayHello((req as HelloRequest).getOrThrow())
}
return v.f(inhr, info, UnaryHandler(handlercs))
case None =>
println("--------------------------111")
println("inhr.name = ${inhr.name}")
println("--------------------------111")
return (srv as (GreeterServer)).getOrThrow().SayHello(inhr)
}
}
public func defaultAnyAndServerStreamToGrpcError(any: Any, ss: ServerStream): GrpcError {
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------0")
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------0")
if (any is HelloRequest) {
let hy: HelloRequest = (any as HelloRequest).getOrThrow()
println("hy.name = ${hy.name}")
} else if (any is HelloReply) {
let hy: HelloReply = (any as HelloReply).getOrThrow()
println("hy.message = ${hy.message}")
}
let hr: HelloRequest = HelloRequest()
hr.name = "测试processStreamingRPC--20221008!仓颉666"
ss.sendMsg(hr)
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------1")
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------1")
return NULL_ERR
}
var Greeter_ServiceDesc = ServiceDesc(
serviceName: "helloworld.Greeter",
handlerType: -1,
methods: [MethodDesc(methodName: "SayHello", handler: Greeter_SayHello_Handler)],
//streams: [StreamDesc(streamName: "SayHello", handler: StreamHandler(defaultAnyAndServerStreamToGrpcError) )],
metadata: "./helloworld.proto"
)
func registerGreeterServer(s: Server, srv: MessageLite) {
s.registerService(Greeter_ServiceDesc, srv)
}
执行结果如下:
--------------------------111
inhr.name = World
--------------------------111
test pass
2 客户端流式RPC交互方式客户端与服务端的信息通信功能
前置条件:NA
场景:
1.提供客户端流式RPC:客户端用提供的一个数据流写入并发送一系列消息给服务端;
2.服务端流式RPC:客户端发送一个请求给服务端,可获取一个数据流用来读取一系列消息。
约束: NA
性能: 支持版本几何性能持平
可靠性: NA
2.1 客户端流式请求发送
2.1.1 主要接口
public class ClientConn {
/*
* 创建流式客户端
* 返回值 ClientStream - 返回一个流式客户端
*/
public func newStreamClient(): ClientStream
}
public class ClientStream <: GrpcStream {
/*
* 指定要请求的服务
* 参数 method - 要请求的服务
*/
public func doStreamRequest(method: String): Unit
/*
* 发送流数据
* 参数 message - 要发送的 protobuf 请求对象
*/
public func sendStreamMsg(message: MessageLite): Unit
/*
* 停止发送并开始接收服务器响应
* 返回值 ReadStream - 返回一个 ReadStream 用于读取接收到的响应数据
*/
public func closeAndRecv(): ReadStream
}
public class ReadStream {
/*
* 从返回流读取 protobuf 对象
* 参数 output - 用于接收读取到的 protobuf 对象
* 返回值 Bool - 是否成功读取到,true 为读取成功, false 为未读到数据
*/
public func recv(output: MessageLite): Bool
}
2.2 服务端流式请求处理
服务端流式请求,在用例中书写注册配置ServiceDesc时,去掉一元配置methods字段,增加流式配置streams字段。
2.2.1 主要接口
2.2.1.1 class ServerStream
实现服务器端流RPC。
public class ServerStream <: Stream {
/**
* 根据可选择参数,创建 ServerStream 实例。
* 参数 stss - ServerTransport 数据结构,是所有gRPC服务器端传输实现的通用父类
* 参数 stream - Stream 数据结构,表示传输层中的RPC
* 参数 parser - Parser 数据结构,解析器从底层读取器读取完整的gRPC消息
* 参数 codec - BaseCodec 数据类型,包含编解码器和编码的功能
* 参数 maxReceiveMessageSize - Int64 数据类型
* 参数 maxSendMessageSize - Int64 数据结构
* 参数 statsHandler - Array<GrpcHandler> 数据类型,GrpcHandler 定义相关统计处理的接口
*/
public init(stss!: ServerTransport,
stream!: Stream = Stream(isNone: true),
parser!: Parser = Parser(),
codec!: BaseCodec = BaseCodec(),
maxReceiveMessageSize!: Int64 = 0,
maxSendMessageSize!: Int64 = 0,
statsHandler!: Array<GrpcHandler> = Array<GrpcHandler>()
)
/**
* 设置头元数据。它可以被多次调用。
* 多次调用时,将合并所有提供的元数据。
* 发生以下情况之一时,将发送所有元数据:
* -调用 sendHeader();
* -发出第一个响应;
* -发送RPC状态(错误或成功)。
* 参数 md - HashMap 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public override func setHeader(md: HashMap<String, Array<String>>): GrpcError
/**
* 发送标头元数据。将发送由 setHeader() 设置的提供的md和标头。
* 如果多次调用,则会失败。
* 参数 md - HashMap 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public override func sendHeader(md: HashMap<String, Array<String>>): GrpcError
/**
* 设置将随RPC状态一起发送的尾部元数据。多次调用时,将合并所有提供的元数据。
* 参数 md - HashMap 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
*/
public override func setTrailer(md: HashMap<String, Array<String>>): GrpcError
/**
* 发送消息。出现错误时,sendMsg中止流,并直接返回错误。
* sendMsg阻止,直到:
* -有足够的流量控制来安排m的运输,或
* -流完成,或
* -溪流中断。
* 不会等到客户端收到消息。过早关闭流可能会导致消息丢失。
* 让一个goroutine同时调用SendMsg和另一个goroutine同时在同一个流上调用RecVMS是安全的,
* 但在不同的goroutines中在同一流上调用SendMsg是不安全的。
* 参数 m - MessageLite 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
* 注意:参数 m 在当前版本支持的参数类型为MessageLite,后续版本更迭会修改参数类型为Any
*/
public func sendMsg(m: MessageLite): GrpcError
/**
* 阻塞,直到它接收到m中的消息或流完成。
* 当客户端执行CloseSend时,它返回io EOF。对于任何非EOF错误,流将被中止,错误包含RPC状态。
* 让一个goroutine同时调用SendMsg和另一个goroutine同时在同一个流上调用RecVMs是安全的,
* 但在不同的goroutines中在同一流上调用RECVMs是不安全的。
* 参数 m - MessageLite 数据结构
* 返回值 GrpcError - GrpcError 数据结构,封装错误信息的类
* 注意:参数 m 在当前版本支持的参数类型为MessageLite,后续版本更迭会修改参数类型为Any
*/
public func recvMsg(m: MessageLite): GrpcError
/**
* 返回当前的 Context
* 返回值 Context - 当前的 Context
*/
public func context():Context
}
2.2.1.2 class ServiceDesc
表示流RPC服务的规范.
public class ServiceDesc {
/**
* 根据可选择参数,创建 ServiceDesc 实例。
* 参数 ctx - Context 数据结构
* 参数 serviceName - String 数据结构
* 参数 handlerType - Any 数据结构, 指向服务接口的指针。
* 参数 methods - Array<MethodDesc> 数据结构
* 参数 streams - Array<StreamDesc> 数据结构
* 参数 metadata - Any 数据结构
*/
public init(ctx: Context,
serviceName!: String = "",
handlerType!: Any,
methods!: Array<MethodDesc> = Array<MethodDesc>(),
streams!: Array<StreamDesc> = Array<StreamDesc>(),
metadata!: Any)
}
2.2.1.3 class StreamHandler
定义gRPC服务器调用的处理程序,以完成流式RPC的执行。
public class StreamHandler {
/**
* 根据参数函数f,创建 StreamHandler 实例。
* 参数 f - 函数类型
*/
public init(f: (Any, ServerStream) -> GrpcError)
}
2.2.1.4 class StreamServerInterceptor
提供了一个钩子来拦截服务器上流式RPC的执行。 info包含拦截器可以操作的RPC的所有信息。处理程序是服务方法实现。拦截器负责调用处理程序来完成RPC。
public class StreamServerInterceptor {
/**
* 根据参数函数f,创建 StreamServerInterceptor 实例。
* 参数 f - 函数类型
*/
public init(f: (Any, ServerStream, StreamServerInfo, StreamHandler) -> GrpcError)
}
2.2.1.5 class StreamServerInfo
包含关于服务器端流式RPC的各种信息。拦截器可以改变所有每rpc信息。
public class StreamServerInfo {
/**
* 根据可选择参数,创建 StreamServerInfo 实例。
* 参数 fullMethod - String 数据结构,完整的RPC方法字符串,即/package.service/method。
* 参数 isClientStream - Bool 数据结构,指示RPC是否是客户端流RPC。
* 参数 isServerStream - Bool 数据结构,指示RPC是否是服务器流RPC。
*/
public init(fullMethod!: String = "",
isClientStream!: Bool = false,
isServerStream!: Bool = false)
}
2.3 示例
from std import socket.*
from rpc4cj import grpc.*
from rpc4cj import transport.*
from rpc4cj import exceptions.*
from rpc4cj import util.*
from std import collection.*
from std import time.*
from std import sync.*
from protobuf import protobuf.*
let strServer: String = "20221226abc"
var port: UInt16 = 0
let num: Int64 = 5
main() {
let ss: SocketServer = SocketServer(TCP, "127.0.0.1", port)
port = ss.port
spawn{ =>
let lis: GRPCServerSocket = GRPCServerSocket(GrpcSocketS(ss))
let server: Server = Server()
registerGreeterServer(server, ServerTest())
server.serve(lis)
}
sleep(Duration.second)
var req = HelloRequest()
req.name = "World"
var con = dial("127.0.0.1:${port}")//linux
let stream = con.newStreamClient()
stream.doStreamRequest("/helloworld.Greeter/SayHello")
let fut: Future<String> = spawn { =>
var resp = HelloReply()
while(stream.recv(resp)) {
println("client-resp = ${resp}")
}
return "finish"
}
for (i in 0..num) {
req.name = "World-" + i.toString()
stream.sendStreamMsg(req)
}
stream.closeSend()
let res: Option<String> = fut.get(1000 * 1000 * 1000 * 1000 )
match (res) {
case Some(v) => println(v)
case None => throw Exception("gprc request timeout")
}
return 0
}
//protobuf 依赖
class ServerTest <: GreeterServer {
public func SayHello(hr: HelloRequest): (HelloReply, GrpcError) {
let hy: HelloReply = HelloReply()
hy.message = "ceshi-20221227"
println("~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~0")
println("received: hr.name = ${hr.name}")
println("~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~1")
return (hy, GrpcError())
}
public func mustEmbedUnimplementedGreeterServer() {}
}
public class HelloRequest <: MessageLite & TypedMessage<HelloRequest> {
private var p_name = Pb3String()
public init() {
m_innerInit()
}
public mut prop name: String { get() { p_name.get() } set(i_vtmp) { p_name.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_name.isEmpty()) { i_tmp += 1 + binSize(p_name.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_name.isEmpty()
}
protected func m_clear(): Unit {
p_name.clear()
}
public func copyFrom(src: HelloRequest): HelloRequest {
p_name = src.p_name
m_innerCopyFrom(src)
return this
}
public func clone(): HelloRequest {
return HelloRequest().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_name.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_name.isEmpty()) { out.append(10); out.packBin(p_name.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_name.isEmpty()) { i_tmp.append("name", p_name.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloRequest() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
}
public class HelloReply <: MessageLite & TypedMessage<HelloReply> {
private var p_message = Pb3String()
public init() {
m_innerInit()
}
public mut prop message: String { get() { p_message.get() } set(i_vtmp) { p_message.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_message.isEmpty()) { i_tmp += 1 + binSize(p_message.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_message.isEmpty()
}
protected func m_clear(): Unit {
p_message.clear()
}
public func copyFrom(src: HelloReply): HelloReply {
p_message = src.p_message
m_innerCopyFrom(src)
return this
}
public func clone(): HelloReply {
return HelloReply().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_message.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_message.isEmpty()) { out.append(10); out.packBin(p_message.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_message.isEmpty()) { i_tmp.append("message", p_message.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloReply() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
}
//server 端
public abstract class GreeterServer <: MessageLite {
public func SayHello(hr: HelloRequest): (HelloReply, GrpcError)
public func mustEmbedUnimplementedGreeterServer(): Unit
protected func m_size(): Int64 { return 0 }
protected func m_isEmpty(): Bool { return false }
protected func m_clear(): Unit {}
public func unpacki(_: BytesReader): Unit {}
public func packi(_: BytesWriter): Unit {}
public func toString(): String { return "" }
}
class UnimplementedGreeterServer <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
return (HelloReply(), GrpcError(Unimplemented, "method SayHello not implemented"))
}
public func mustEmbedUnimplementedGreeterServer() {}
}
interface UnsafeGreeterServer {
func mustEmbedUnimplementedGreeterServer(): Unit
}
public func defaultAnyAndServerStreamToGrpcError(any: Any, ss: ServerStream): GrpcError {
if (any is HelloRequest) {
let hy: HelloRequest = (any as HelloRequest).getOrThrow()
println("hy.name = ${hy.name}")
} else if (any is HelloReply) {
let hy: HelloReply = (any as HelloReply).getOrThrow()
println("hy.message = ${hy.message}")
}
let hr: HelloRequest = HelloRequest()
for (i in 0..num) {
hr.name = "ceshi-20221227-stream-" + i.toString()
ss.sendMsg(hr)
}
var resp = HelloReply()
for (_ in 0..num) {
ss.recvMsg(resp)
println("server-resp = ${resp}")
}
return NULL_ERR
}
var Greeter_ServiceDesc = ServiceDesc(
serviceName: "helloworld.Greeter",
handlerType: -1,
streams: [StreamDesc(streamName: "SayHello", handler: StreamHandler(defaultAnyAndServerStreamToGrpcError) )],
metadata: "./helloworld.proto"
)
func registerGreeterServer(s: Server, srv: MessageLite) {
s.registerService(Greeter_ServiceDesc, srv)
}
执行结果如下:
server-resp = {message: "World-0"}
server-resp = {message: "World-1"}
server-resp = {message: "World-2"}
server-resp = {message: "World-3"}
server-resp = {message: "World-4"}
client-resp = {message: "ceshi-20221227-stream-0"}
client-resp = {message: "ceshi-20221227-stream-1"}
client-resp = {message: "ceshi-20221227-stream-2"}
client-resp = {message: "ceshi-20221227-stream-3"}
client-resp = {message: "ceshi-20221227-stream-4"}
3 提供双向流式 RPC交互方式客户端与服务端的信息通信功能
前置条件:NA
场景:
1.提供双向流式 RP交互方式:客户端和服务端都可以分别通过一个读写数据流来发送一系列消息;
2.提供客户端数据请求数据压缩处理;
3.服务端数据返回数据压缩处理;
约束:压缩解压缩功能与其他语言不通用
性能: 支持版本几何性能持平
可靠性: NA
3.1 客户端数据请求数据压缩处理
3.1.1 主要接口
public class ClientConn {
/*
* 启用默认 zlib 内置压缩器压缩请求数据
*/
public func useCompressor(): Unit
}
3.2 服务端数据返回数据压缩处理
服务端双向流式数据交互,服务端数据返回数据压缩处理,不增加新的api接口。 根据 Client 端的请求头字段 grpc-encoding 判断是否进行数据压缩解压缩处理。 压缩解压缩功能与其他语言不通用。
3.3 示例
from std import socket.*
from rpc4cj import grpc.*
from rpc4cj import transport.*
from rpc4cj import exceptions.*
from rpc4cj import util.*
from std import time.*
from std import sync.*
from std import collection.*
from protobuf import protobuf.*
let strServer: String = "test pass"
class ServerTest <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
let hy: HelloReply = HelloReply()
hy.message = strServer
return (hy, GrpcError())
}
public func mustEmbedUnimplementedGreeterServer() {}
}
var port: UInt16 = 50051
main() {
let ss: SocketServer = SocketServer(TCP, "127.0.0.1", port)
port = ss.port
spawn{ =>
let lis: GRPCServerSocket = GRPCServerSocket(GrpcSocketS(ss))
let server: Server = Server()
registerGreeterServer(server, ServerTest())
server.serve(lis)
}
sleep(Duration.second)
var req = HelloRequest()
var resp = HelloReply()
req.name = "World"
var con = dial("127.0.0.1:${port}")//linux
con.useCompressor()
con.invoke("/helloworld.Greeter/SayHello",req,resp)
println(resp)
return 0
}
public class HelloRequest <: MessageLite & TypedMessage<HelloRequest> {
private var p_name = Pb3String()
public init() {
m_innerInit()
}
public mut prop name: String { get() { p_name.get() } set(i_vtmp) { p_name.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_name.isEmpty()) { i_tmp += 1 + binSize(p_name.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_name.isEmpty()
}
protected func m_clear(): Unit {
p_name.clear()
}
public func copyFrom(src: HelloRequest): HelloRequest {
p_name = src.p_name
m_innerCopyFrom(src)
return this
}
public func clone(): HelloRequest {
return HelloRequest().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_name.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_name.isEmpty()) { out.append(10); out.packBin(p_name.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_name.isEmpty()) { i_tmp.append("name", p_name.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloRequest() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
}
public class HelloReply <: MessageLite & TypedMessage<HelloReply> {
private var p_message = Pb3String()
public init() {
m_innerInit()
}
public mut prop message: String { get() { p_message.get() } set(i_vtmp) { p_message.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_message.isEmpty()) { i_tmp += 1 + binSize(p_message.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_message.isEmpty()
}
protected func m_clear(): Unit {
p_message.clear()
}
public func copyFrom(src: HelloReply): HelloReply {
p_message = src.p_message
m_innerCopyFrom(src)
return this
}
public func clone(): HelloReply {
return HelloReply().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_message.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_message.isEmpty()) { out.append(10); out.packBin(p_message.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_message.isEmpty()) { i_tmp.append("message", p_message.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloReply() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
}
//server 端
public abstract class GreeterServer <: MessageLite {
public func SayHello(hr: HelloRequest): (HelloReply, GrpcError)
public func mustEmbedUnimplementedGreeterServer(): Unit
protected func m_size(): Int64 { return 0 }
protected func m_isEmpty(): Bool { return false }
protected func m_clear(): Unit {}
public func unpacki(_: BytesReader): Unit {}
public func packi(_: BytesWriter): Unit {}
public func toString(): String { return "" }
}
class UnimplementedGreeterServer <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
return (HelloReply(), GrpcError(Unimplemented, "method SayHello not implemented"))
}
public func mustEmbedUnimplementedGreeterServer() {}
}
interface UnsafeGreeterServer {
func mustEmbedUnimplementedGreeterServer(): Unit
}
public func defaultAnyAndServerStreamToGrpcError(any: Any, ss: ServerStream): GrpcError {
if (any is HelloRequest) {
let hy: HelloRequest = (any as HelloRequest).getOrThrow()
println("hy.name = ${hy.name}")
} else if (any is HelloReply) {
let hy: HelloReply = (any as HelloReply).getOrThrow()
println("hy.message = ${hy.message}")
}
let hr: HelloRequest = HelloRequest()
hr.name = "测试processStreamingRPC--20221008!仓颉666"
ss.sendMsg(hr)
return NULL_ERR
}
var Greeter_ServiceDesc = ServiceDesc(
serviceName: "helloworld.Greeter",
handlerType: -1,
//methods: [MethodDesc(methodName: "SayHello", handler: Greeter_SayHello_Handler)],
streams: [StreamDesc(streamName: "SayHello", handler: StreamHandler(defaultAnyAndServerStreamToGrpcError) )],
metadata: "./helloworld.proto"
)
func registerGreeterServer(s: Server, srv: MessageLite) {
s.registerService(Greeter_ServiceDesc, srv)
}
执行结果如下:
{message: "测试processStreamingRPC--20221008!仓颉666"}
4 提供gPPC配置功能
前置条件:NA
场景:
1.提供服务端请求拦截器;
2.服务端请求拦截器;
3.服务器启动选项控制。
约束:支持的拦截包括对请求方法名,请求信息和返回值信息的拦截
性能: 支持版本几何性能持平
可靠性: NA
4.1 客户端请求拦截器
实现拦截器接口,调用对应函数,对客户设置拦截器。
4.1.1 主要接口
public interface UnaryInvoker {
/*
* 拦截一元 inkove 请求
* 参数 method - 请求的 grpc 方法名
* 参数 req - 请求的参数
* 参数 resp - 请求的返回值
*/
func orderUnaryClientInterceptor(method: String, req: MessageLite, resp: MessageLite): Unit
}
public interface StreamClientInterceptor {
/*
* 拦截客户端流式请求方法名
* 参数 method - 请求的 grpc 方法名
*/
func orderStreamClientInterceptor(method: String): Unit
/*
* 拦截客户端流式请求参数
* 参数 req - 请求的参数
*/
func orderStreamClientInterceptorReq(req: MessageLite): Unit
/*
* 拦截客户端流式请求返回值
* 参数 resp - 请求的返回值
*/
func orderStreamClientInterceptorResp(resp: MessageLite): Unit
}
4.2 服务端请求拦截器
调用对应函数,对服务器选项设置流拦截器,一元拦截器。
4.2.1 主要接口
/**
* 返回一个 ServerOption,该选项指定流式RPC的链接拦截器。
* 第一个拦截器将是最外层的,而最后一个拦截程序将是真正调用的最内层包装器。
* 通过此方法添加的所有流拦截器将被链接。
* 参数 interceptors - Array<StreamServerInterceptor> 数据类型
* 返回值 ServerOption - 接口类型
*/
public func chainStreamInterceptor(interceptors: Array<StreamServerInterceptor>): ServerOption
/**
* 返回一个 ServerOption,该选项指定一元RPC的链接拦截器。
* 第一个拦截器将是最外层的,而最后一个拦截程序将是真正调用的最内层包装器。
* 通过此方法添加的所有一元拦截器都将被链接。
* 参数 interceptors - Array<UnaryServerInterceptor> 数据类型
* 返回值 ServerOption - 接口类型
*/
public func chainUnaryInterceptor(interceptors: Array<UnaryServerInterceptor>): ServerOption
/**
* 返回一个 ServerOption,该选项设置服务器的StreamServerInterceptor。
* 只能安装一个流拦截器。
* 参数 si - StreamServerInterceptor 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func streamInterceptor(si: StreamServerInterceptor): ServerOption
/**
* 返回设置服务器的UnaryServerInterceptor的ServerOption。
* 只能安装一个一元拦截器。
* 多个拦截器的构造(例如链接)可以在调用方实现。
* 参数 usi - UnaryServerInterceptor 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func unaryInterceptor(usi: UnaryServerInterceptor): ServerOption
4.3 服务器启动选项控制
调用对应函数,对服务器选项设置诸如凭据、编解码器和保留活动参数等选项进行赋值。
4.3.1 主要接口
/**
* 返回一个 ServerOption,该选项设置所有新连接的连接建立超时(包括HTTP/2握手)。
* 如果未设置,默认值为120秒。
* 零值或负值将导致立即超时。
* 实验的
* 注意:此API是实验性的,可能会在以后的版本中更改或删除。
* 参数 d - Duration 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func connectionTimeout(d: Duration): ServerOption
/**
* 返回一个 ServerOption,该选项设置用于消息封送和解封送的编解码器。
* 这将覆盖通过RegisterCodec注册的编解码器的任何内容子类型查找。
* 不推荐使用:使用encoding.RegisterCodec注册编解码器。服务器将根据传入请求的头自动使用注册的编解码器。
* 参数 codec - BaseCodec 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func customCodec(codec: BaseCodec): ServerOption
/**
* 返回一个 ServerOption,用于设置流的动态头表的大小。
* 参数 num - UInt32 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func headerTableSize(num: UInt32): ServerOption
/**
* 返回一个 ServerOption,该选项设置要创建的所有服务器传输的点击句柄。只能安装一个。
* 参数 handle - ServerInHandle 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func inTapHandle(handle: ServerInHandle): ServerOption
/**
* 返回设置连接窗口大小的服务器选项。
* 窗口大小的下限为64K,任何小于该值的值都将被忽略。
* 参数 s - Int32 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func initialConnWindowSize(s: Int32): ServerOption
/**
* 返回为服务器设置keepalive强制策略的ServerOption。
* 参数 ep - EnforcementPolicy 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func keepaliveEnforcementPolicy(ep: EnforcementPolicy): ServerOption
/**
* 返回一个 ServerOption,用于设置服务器的keepalive和max age参数。
* 参数 kp - ServerParameters 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func keepaliveParams(kp: ServerParameters): ServerOption
/**
* 返回一个 ServerOption,该选项将对每个ServerTransport的并发流数量施加限制。
* 参数 mc - UInt32 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func maxConcurrentStreams(mc: UInt32): ServerOption
/**
* 返回一个 ServerOption,该选项设置服务器准备接受的头列表的最大(未压缩)大小。
* 参数 num - UInt32 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func maxHeaderListSize(num: UInt32): ServerOption
/**
* MaxMsgSize返回一个 ServerOption,以设置服务器可以接收的最大消息大小(以字节为单位)。
* 如果未设置,gRPC将使用默认限制。
* 不推荐使用:改用MaxRecvMsgSize。
* 参数 m - Int64 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func maxMsgSize(m: Int64): ServerOption
/**
* 返回服务器选项,以设置服务器可以接收的最大消息大小(以字节为单位)。
* 如果未设置,gRPC将使用默认的4MB。
* 参数 mr - Int64 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func maxRecvMsgSize(mr: Int64): ServerOption
/**
* 返回服务器选项,以设置服务器可以发送的最大消息大小(以字节为单位)。
* 如果未设置,gRPC将使用默认的“math.MaxInt32”。
* 参数 num - Int64 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func maxSendMsgSize(num: Int64): ServerOption
/**
* 返回一个 ServerOption,该选项设置应用于处理传入流的工作进程的数量。
* 将其设置为零(默认值)将禁用worker并为每个流生成一个新的线程。
* 注意:此API是实验性的,可能会在以后的版本中更改或删除。
* 参数 num - UInt32 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func numStreamWorkers(num: UInt32): ServerOption
/**
* 允许您设置读取缓冲区的大小,这决定了一次读取系统调用最多可以读取多少数据。
* 此缓冲区的默认值为32KB。
* Zero将禁用连接的读取缓冲区,以便数据成帧器可以直接访问底层conn。
* 参数 s - Int64 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func readBufferSize(s: Int64): ServerOption
/**
* 返回一个 ServerOption,用于设置服务器的统计处理程序。
* 参数 h - GrpcHandler 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func statsHandler(h: GrpcHandler): ServerOption
/**
* 返回允许添加自定义未知服务处理程序的ServerOption。
* 提供的方法是一个bidi-streaming RPC服务处理程序,每当收到未注册服务或方法的请求时,将调用该处理程序,而不是返回“未实现”的gRPC错误。
* 处理函数和流拦截器(如果设置)可以完全访问ServerStream,包括其上下文。
* 参数 streamHandler - StreamHandler 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func unknownServiceHandler(streamHandler: StreamHandler): ServerOption
/**
* 确定在对线路进行写入之前可以批处理多少数据。
* 该缓冲区的相应内存分配将是保持系统调用低的大小的两倍。
* 此缓冲区的默认值为32KB。
* 零将禁用写入缓冲区,以便每次写入都在基础连接上。
* 注意:发送呼叫可能不会直接转换为写入。
* 参数 s - Int64 数据类型
* 返回值 ServerOption - 接口类型,返回为服务器连接设置凭据的ServerOption。
*/
public func writeBufferSize(s: Int64): ServerOption
4.4 示例
from std import socket.*
from rpc4cj import grpc.*
from rpc4cj import transport.*
from rpc4cj import exceptions.*
from rpc4cj import util.*
from std import time.*
from std import sync.*
from std import collection.*
from protobuf import protobuf.*
let strServer: String = "World"
class ServerTest <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
let hy: HelloReply = HelloReply()
hy.message = strServer
return (hy, GrpcError())
}
public func mustEmbedUnimplementedGreeterServer() {}
}
var port: UInt16 = 50052
main() {
let ss: SocketServer = SocketServer(TCP, "127.0.0.1", port)
port = ss.port
spawn{ =>
let lis: GRPCServerSocket = GRPCServerSocket(GrpcSocketS(ss))
let ssi1: StreamServerInterceptor = getInstance_StreamServerInterceptor()
let so1: ServerOption = streamInterceptor(ssi1)
var opts: Array<ServerOption> = Array<ServerOption>([so1])
let server: Server = Server.newServer(opts)
registerGreeterServer(server, ServerTest())
server.serve(lis)
}
sleep(Duration.second)
var req = HelloRequest()
var resp = HelloReply()
req.name = "World"
var opts = ArrayList<DialOption>()
opts.append(WithUnaryInterceptor(unary()))
var con = dial("127.0.0.1:${port}",opts )//linux
con.invoke("/helloworld.Greeter/SayHello",req,resp)
return 0
}
func funcIns_StreamServerInterceptor(_: Any,_: ServerStream, _: StreamServerInfo ,_: StreamHandler): GrpcError {
println("测试流式拦截器调用!")
return GrpcError()
}
public func getInstance_StreamServerInterceptor(): StreamServerInterceptor{
return StreamServerInterceptor(funcIns_StreamServerInterceptor)
}
class unary <: UnaryInvoker {
public func orderUnaryClientInterceptor(method: String, req: MessageLite, _: MessageLite): Unit {
if (method == "/helloworld.Greeter/SayHello") {
println(method)
println("method ok")
} else {
throw Exception("error method")
}
var r = (req as HelloRequest).getOrThrow()
if (r.name == "World") {
println(r.name)
println("request ok")
} else {
throw Exception("error method")
}
}
}
public class HelloRequest <: MessageLite & TypedMessage<HelloRequest> {
private var p_name = Pb3String()
public init() {
m_innerInit()
}
public mut prop name: String { get() { p_name.get() } set(i_vtmp) { p_name.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_name.isEmpty()) { i_tmp += 1 + binSize(p_name.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_name.isEmpty()
}
protected func m_clear(): Unit {
p_name.clear()
}
public func copyFrom(src: HelloRequest): HelloRequest {
p_name = src.p_name
m_innerCopyFrom(src)
return this
}
public func clone(): HelloRequest {
return HelloRequest().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_name.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_name.isEmpty()) { out.append(10); out.packBin(p_name.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_name.isEmpty()) { i_tmp.append("name", p_name.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloRequest() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
}
public class HelloReply <: MessageLite & TypedMessage<HelloReply> {
private var p_message = Pb3String()
public init() {
m_innerInit()
}
public mut prop message: String { get() { p_message.get() } set(i_vtmp) { p_message.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_message.isEmpty()) { i_tmp += 1 + binSize(p_message.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_message.isEmpty()
}
protected func m_clear(): Unit {
p_message.clear()
}
public func copyFrom(src: HelloReply): HelloReply {
p_message = src.p_message
m_innerCopyFrom(src)
return this
}
public func clone(): HelloReply {
return HelloReply().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_message.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_message.isEmpty()) { out.append(10); out.packBin(p_message.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_message.isEmpty()) { i_tmp.append("message", p_message.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloReply() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
}
//server 端
public abstract class GreeterServer <: MessageLite {
public func SayHello(hr: HelloRequest): (HelloReply, GrpcError)
public func mustEmbedUnimplementedGreeterServer(): Unit
protected func m_size(): Int64 { return 0 }
protected func m_isEmpty(): Bool { return false }
protected func m_clear(): Unit {}
public func unpacki(_: BytesReader): Unit {}
public func packi(_: BytesWriter): Unit {}
public func toString(): String { return "" }
}
class UnimplementedGreeterServer <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
return (HelloReply(), GrpcError(Unimplemented, "method SayHello not implemented"))
}
public func mustEmbedUnimplementedGreeterServer() {}
}
interface UnsafeGreeterServer {
func mustEmbedUnimplementedGreeterServer(): Unit
}
public func defaultAnyAndServerStreamToGrpcError(any: Any, ss: ServerStream): GrpcError {
if (any is HelloRequest) {
let hy: HelloRequest = (any as HelloRequest).getOrThrow()
println("hy.name = ${hy.name}")
} else if (any is HelloReply) {
let hy: HelloReply = (any as HelloReply).getOrThrow()
println("hy.message = ${hy.message}")
}
let hr: HelloRequest = HelloRequest()
hr.name = "测试processStreamingRPC--20221008!仓颉666"
ss.sendMsg(hr)
return NULL_ERR
}
var Greeter_ServiceDesc = ServiceDesc(
serviceName: "helloworld.Greeter",
handlerType: -1,
//methods: [MethodDesc(methodName: "SayHello", handler: Greeter_SayHello_Handler)],
streams: [StreamDesc(streamName: "SayHello", handler: StreamHandler(defaultAnyAndServerStreamToGrpcError) )],
metadata: "./helloworld.proto"
)
func registerGreeterServer(s: Server, srv: MessageLite) {
s.registerService(Greeter_ServiceDesc, srv)
}
执行结果如下:
测试流式拦截器调用!
/helloworld.Greeter/SayHello
method ok
World
request ok
5 提供域名解析功能
前置条件:NA
场景:
1.实现一个内置dns解析器,作为grpc中的默认解析器安装和包传递实现了传递解析程序,把域名转换成为网络可以识别的IP地址;
2.允许自定义了一个解析程序,可用于将解析的地址手动发送到ClientConn,允许用户添加自定义域名器。
约束: NA
性能: 支持版本几何性能持平
可靠性: NA
5.1 域名解析
5.1.1 主要接口
public struct Address {
public var addr: String
public var serverName: String
public var addrType: UInt8
public var attributes: attributes.Attributes<Int64,String>
/**
* 判断是否和另一个 Address 的所有值是否相等
* 参数 o - 要对比的另一个 Address 的所有值
* 返回值 Bool - 值是否相等
*/
public func equal(o: Address): Bool {
}
/**
* 创建一个 Attributes<K,V>
* 参数 key - 一个 key 值,该类型必须 实现了Hashable和Equatable<K>
* 参数 value - 一个 value 值,该类型必须 实现了Equatable<V>
* 返回值 Bool - 值是否相等
*/
public func new<K,V>(key: K, value: V): Attributes<K,V> where K <: Hashable & Equatable<K>, V <: Equatable<V>
/**
* 创建一个空 AddressMap 对象
* 返回值 AddressMap - 返回一个空 AddressMap 对象
*/
public func newAddressMap(): AddressMap
public class AddressMap {
/**
* 设置一个 Address
* 参数 addr - 设置的 Address
* 参数 value - 设置对应的字符串值
*/
public func set(addr: Address, value: String): Unit
/**
* 获取 Address 对应的值
* 参数 addr - 要获取的 Address
* 返回值 (String, Bool) - 值和是否成功获取
*/
public func get(addr: Address): (String, Bool)
/**
* 移除一个 Address
* 参数 addr - 要删除 Address
*/
public func delete(addr: Address): Unit
/**
* 当前存入的 Address 数量
* 返回值 Int64 - 返回大小
*/
public func len(): Int64
/**
* 获取所有的 Address
* 返回值 ArrayList<Address> - 返回 Address 列表
*/
public func keys(): ArrayList<Address>
}
/**
* 设置默认解析头
* 参数 scheme - 设置默认解析头字符串
*/
public func setDefaultScheme(scheme: String)
/**
* 获取默认解析头
* 返回值 String - 返回当前解析头
*/
public func getDefaultScheme(): String
/**
* 将一个构造器 Builder 注册进系统
* 参数 b - 要注册的构造器
*/
public func register(b: Builder): Unit
/**
* 根据解析获取一个构造器
* 参数 scheme - 解析头字符串
* 返回值 Builder - 返回对应的构造器,没有获取到则内部抛出异常
*/
public func get(scheme: String): Builder
/**
* 根据解析头取消构造器的注册
* 参数 scheme - 解析头字符串
*/
public func unregisterForTesting(scheme: String): Unit
/*
* 配置解析器
* 参数 rs - 配置的解析器
* 返回值 DialOption - 返回的 DialOption
*/
public func withResolvers(rs: resolver.Builder): DialOption
/**
* 获取内置的 dns 解析器 Builder 用于设置域名解析。
* 返回值 resolver.Builder - 返回一个 dns 解析器
*/
public func newDnsBuilder(): resolver.Builder
/**
* 通过 scheme 名创建一个空 Resolver 对象。
* 返回值 Resolver - 返回一个空 Resolver 对象
*/
public func newBuilderWithScheme(scheme: String): Resolver
public open class ManualResolver <: resolver.Builder & resolver.Resolver {
/**
* 初始化状态
* 参数 scheme - 解析头字符串
*/
public func initialState(s: resolver.State)
/**
* 传入解析的目标,构造解析器
* 参数 target - 解析目标
* 返回值 resolver.Resolver - 返回对应的解析器
* 返回值 GrpcError - 返回错误信息
*/
public func build(target: resolver.Target): (resolver.Resolver, GrpcError)
/**
* 返回当前解析器的头
* 参数 String - 返回解析头字符串
*/
public func scheme(): String
/**
* 用户传入自定义参数做自定义解析
* 参数 o - ResolveNowOptions 类型的自定义参数
*/
public func resolveNow(o: resolver.ResolveNowOptions): Unit
/**
* 调用用户自定义的 ClientConn 的 close
*/
public func close(): Unit
/**
* 调用用户自定义的 ClientConn 的 updateState
* 参数 resolver.State - 要处理的 State
* 返回值 GrpcError - 返回一个 GrpcError
*/
public func updateState(s: resolver.State): GrpcError
/**
* 处理 grpc 错误,调用用户自定义的 ClientConn 的 reportError
* 参数 err - 要处理的 GrpcError
*/
public func reportError(err: GrpcError): Unit
}
/**
* 设置 passthrough 解析器为默认解析器
*/
public func initPassThrough()
/**
* 获取一个默认 Builder 解析器
* 返回值 resolver.Builder - 返回一个 Builder 解析器
*/
public func NewBuilder(): resolver.Builder
5.2 示例
from std import socket.*
from rpc4cj import grpc.*
from rpc4cj import transport.*
from rpc4cj import exceptions.*
from rpc4cj import util.*
from std import time.*
from std import sync.*
from std import collection.*
from rpc4cj import resolver.*
from rpc4cj import resolver.dns.*
var port: UInt16 = 0
main() {
let ss: SocketServer = SocketServer(TCP, "127.0.0.1", port)
port = ss.port
spawn{ =>
let lis: GRPCServerSocket = GRPCServerSocket(GrpcSocketS(ss))
let server: Server = Server()
server.serve(lis)
}
sleep(Duration.second)
try{
dial("passthrough:///127.0.0.1:${port}")//reslove localhost with
} catch (e: Exception) {
println("reslove failed")
return 1
}
println("reslove pass")
return 0
}
执行结果如下:
reslove pass
6 支持TLS安全协议
前置条件:NA
场景:
1.TLS/mTLS,单向tls及双向tls, 是网络传输层安全协议,为在为两个 grpc应用程序之间的通信提供隐私性和数据完整性。用于在 grpc客户端应用程序和服务器端应用程序之间提供安全连接。
约束: tls安全协议只支持仓颉语言
性能: 支持版本几何性能持平
可靠性: NA
6.1 支持TLS安全协议
TLS安全协议,不增加新的api接口。Client 端根据 dial 函数的参数 cfg 判断是否进行使用 tls 安全协议。Server 端在用例书写时可自行确定是进行普通 TCP 连接还是使用 tls。其他语言不通用。
6.1 示例
// DEPENDENCE: ./example-cert.pem
// DEPENDENCE: ./example-key.pem
// EXEC: cjc %import-path %L %l %f %project-path %project-L/protobuf -l protobuf-cj
// EXEC: ./main
from std import socket.*
from std import io.*
from std import time.*
from std import sync.*
from std import collection.*
from net import tls.*
from rpc4cj import grpc.*
from rpc4cj import transport.*
from rpc4cj import exceptions.*
from rpc4cj import util.*
from protobuf import protobuf.*
class ServerTest <: GreeterServer {
// SayHello implements helloworld.GreeterServer
public func SayHello(hr: HelloRequest): (HelloReply, GrpcError) {
let hy: HelloReply = HelloReply()
hy.message = "测试server端-20230224......哈哈哈哈,测试通了!"
println("~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~0")
println("~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~1")
println("测试server端-20220919")
println("received: hr.name = ${hr.name}")
println("~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~2")
println("~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~3")
return (hy, GrpcError())
}
public func mustEmbedUnimplementedGreeterServer() {}
}
var port: UInt16 = 50051
main(): Int64 {
let pem = String.fromUtf8(readToEnd(File("./example-cert.pem", OpenMode.Read)))
let key = String.fromUtf8(readToEnd(File("./example-key.pem", OpenMode.Read)))
var cfgs = TlsServerConfig(X509Certificate.decodeFromPem(pem),GeneralPrivateKey.decodeFromPem(key))
cfgs.supportedAlpnProtocols = ["h2", "http/1.1"]
let ss: TlsSocketServer = TlsSocketServer(TCP, "127.0.0.1", port, cfgs)
port = ss.port
spawn{ =>
println("ss=${ss.address}")//127.0.0.1:50051/hello
let lis: GRPCServerSocket = GRPCServerSocket(GrpcTlsS(ss))
let server: Server = Server()
registerGreeterServer(server, ServerTest())
server.serve(lis)
}
sleep(Duration.second)
var req = HelloRequest()
var resp = HelloReply()
req.name = "World"
var cfg = TlsClientConfig("./example-cert.pem")
cfg.alpnList = ["h2", "http/1.1"]
//var cfg = TlsClientConfig("")
cfg.verifyMode = VerifyNone
try {
//var con = dial("192.168.122.1:${port}")//windows
var con = dial("127.0.0.1:${port}", cfg: cfg)//linux
con.invoke("/helloworld.Greeter/SayHello", req, resp)
} catch(e: Exception) {
e.printStackTrace()
return -1
}
println("result: ${resp}")
println("result: {\"message\": \"Hello World\"}")
return 0
}
public class HelloRequest <: MessageLite & TypedMessage<HelloRequest> {
private var p_name = Pb3String()
public init() {
m_innerInit()
}
public mut prop name: String { get() { p_name.get() } set(i_vtmp) { p_name.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_name.isEmpty()) { i_tmp += 1 + binSize(p_name.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_name.isEmpty()
}
protected func m_clear(): Unit {
p_name.clear()
}
public func copyFrom(src: HelloRequest): HelloRequest {
p_name = src.p_name
m_innerCopyFrom(src)
return this
}
public func clone(): HelloRequest {
return HelloRequest().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_name.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_name.isEmpty()) { out.append(10); out.packBin(p_name.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_name.isEmpty()) { i_tmp.append("name", p_name.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloRequest() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloRequest(); i_tmp.unpack(src); i_tmp }
}
public class HelloReply <: MessageLite & TypedMessage<HelloReply> {
private var p_message = Pb3String()
public init() {
m_innerInit()
}
public mut prop message: String { get() { p_message.get() } set(i_vtmp) { p_message.set(i_vtmp); markDirty() } }
protected func m_size(): Int64 {
var i_tmp = 0
if (!p_message.isEmpty()) { i_tmp += 1 + binSize(p_message.data) }
return i_tmp
}
protected func m_isEmpty(): Bool {
p_message.isEmpty()
}
protected func m_clear(): Unit {
p_message.clear()
}
public func copyFrom(src: HelloReply): HelloReply {
p_message = src.p_message
m_innerCopyFrom(src)
return this
}
public func clone(): HelloReply {
return HelloReply().copyFrom(this)
}
public func unpack<T>(src: T): Unit where T <: BytesReader {
while (!src.isEmpty()) {
let i_tmp = src.parseTag()
match (i_tmp[0]) {
case 1 => p_message.unpack(src)
case _ => m_unknown(src, i_tmp)
}
}
markDirty()
}
public func pack<T>(out: T): Unit where T <: BytesWriter {
if (!p_message.isEmpty()) { out.append(10); out.packBin(p_message.data) }
m_unknown(out)
}
public func toString(): String {
var i_tmp = StructPrinter()
if (!p_message.isEmpty()) { i_tmp.append("message", p_message.get()) }
return i_tmp.done().toString()
}
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
public static func empty() { HelloReply() }
public static func fromBytes(src: Collection<Byte>) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
public static func fromBytes(src: BytesReader) { let i_tmp = HelloReply(); i_tmp.unpack(src); i_tmp }
}
//server端
public abstract class GreeterServer <: MessageLite {
public func SayHello(hr: HelloRequest): (HelloReply, GrpcError)
public func mustEmbedUnimplementedGreeterServer(): Unit
protected func m_size(): Int64 { return 0 }
protected func m_isEmpty(): Bool { return false }
protected func m_clear(): Unit {}
public func unpacki(_: BytesReader): Unit {}
public func packi(_: BytesWriter): Unit {}
public func toString(): String { return "" }
}
class UnimplementedGreeterServer <: GreeterServer {
public func SayHello(_: HelloRequest): (HelloReply, GrpcError) {
return (HelloReply(), GrpcError(Unimplemented, "method SayHello not implemented"))
}
//public func clone() { this }
public func mustEmbedUnimplementedGreeterServer() {}
}
interface UnsafeGreeterServer {
func mustEmbedUnimplementedGreeterServer(): Unit
}
/*
public class DEFAULT_MESSAGELITE2 <: MessageLite {
public func clone() { this }
}
*/
func Greeter_SayHello_Handler(srv: MessageLite, dec: (MessageLite) -> GrpcError, interceptor: ?UnaryServerInterceptor): (MessageLite, GrpcError) {
let inhr: HelloRequest = HelloRequest()
var err: GrpcError = dec(inhr)
if (!err.isNull()) {
return (inhr, err)
}
match(interceptor){
case Some(v) =>
let info: UnaryServerInfo = UnaryServerInfo(
server: srv,
fullMethod: "/helloworld.Greeter/SayHello"
)
println("--------------------------000")
println("inhr.name = ${inhr.name}")
println("--------------------------000")
inhr.name = "dec(inhr)解析到的HelloRequest是空的???-----000"
func handlercs(req: Any): (Any, GrpcError) {
return (srv as (GreeterServer)).getOrThrow().SayHello((req as HelloRequest).getOrThrow())
//return srv.SayHello((req as HelloRequest).getOrThrow())
}
return v.f(inhr, info, UnaryHandler(handlercs))
case None =>
println("--------------------------111")
println("inhr.name = ${inhr.name}")
println("--------------------------111")
//inhr.name = "dec(inhr)解析到的HelloRequest是空的???-----111"
return (srv as (GreeterServer)).getOrThrow().SayHello(inhr)
//return srv.SayHello(inhr)
}
}
public func defaultAnyAndServerStreamToGrpcError(any: Any, ss: ServerStream): GrpcError {
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------0")
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------0")
if (any is HelloRequest) {
let hy: HelloRequest = (any as HelloRequest).getOrThrow()
println("hy.name = ${hy.name}")
} else if (any is HelloReply) {
let hy: HelloReply = (any as HelloReply).getOrThrow()
println("hy.message = ${hy.message}")
}
let hr: HelloRequest = HelloRequest()
hr.name = "测试processStreamingRPC--20221008!仓颉666"
ss.sendMsg(hr)
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------1")
println("++++++++++++++++++++++++++start-------------------defaultAnyAndServerStreamToGrpcError-------1")
return NULL_ERR
}
var Greeter_ServiceDesc = ServiceDesc(
serviceName: "helloworld.Greeter",
handlerType: -1,
methods: [MethodDesc(methodName: "SayHello", handler: Greeter_SayHello_Handler)],
//streams: [StreamDesc(streamName: "SayHello", handler: StreamHandler(defaultAnyAndServerStreamToGrpcError) )],
metadata: "./helloworld.proto"
)
func registerGreeterServer(s: Server, srv: MessageLite) {
s.registerService(Greeter_ServiceDesc, srv)
}
执行结果如下:
ss=127.0.0.1:50051
--------------------------111
inhr.name = World
--------------------------111
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~0
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~1
测试server端-20220919
received: hr.name = World
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~2
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~3
result: {message: "测试server端-20230224......哈哈哈哈,测试通了!"}
result: {"message": "Hello World"}
其他接口
介绍
这些接口在内部流程中使用, 无需调用者操作, 也有可能随迭代修改其实现,请按需谨慎使用。
接口
Http2Client
/**
* grpc 客户端内部请求发送处理类
*/
public class Http2Client {
/*
* 通过 H2ClientConn 发起 grpc 请求,内部接口
* 参数 target - 目标服务器地址
* 参数 h2conn - HTTP2 客户端
* 参数 hdr - grpc 请求头
* 参数 data - grpc 请求体
* 返回值 Array<UInt8> - 收到的原始响应数据
*/
public func write(target: String, h2conn: H2ClientConn, hdr: Array<Byte>, data: Array<Byte>): Array<UInt8>
/*
* 从服务器响应中获取原始响应数据,内部接口
* 参数 resp - HTTP2 服务器响应的 H2Response 对象
* 返回值 ArrayList<Array<UInt8>> - 收到的原始响应数据
*/
public func handleResponse(resp: H2Response): Array<UInt8>
/*
* 从服务器流式响应中获取原始响应数据,内部接口
* 参数 resp - HTTP2 服务器响应的流式 H2Response 对象
* 返回值 ArrayList<Array<UInt8>> - 收到的原始响应数据
*/
public func handleStreamResponse(resp: H2Response): ArrayList<Array<UInt8>>
/*
* 处理 grpc 错误代码,根据错误码抛出响应异常,内部接口
* 参数 errCode - grpc 标准错误代码字符串
*/
public func handleError(errCode: String): Unit
}
class Attributes
/**
* grpc 客户端内部键值对存储类型
*/
public class Attributes<K,V> where K <: Hashable & Equatable<K>, V <: Equatable<V> {
/*
* 判断是否和当前 Attributes 相等
* 参数 o - 要对比的 Attributes
* 返回值 Bool - 是否相等
*/
public func Equal(o: Attributes<K,V>): Bool
}
/*
* 根据键和值创建一个 Attributes 对象
* 参数 key - 要传入的键
* 参数 value - 要传入的值
* 返回值 Attributes<K,V> - 返回一个 Attributes 对象
*/
public func new<K,V>(key: K, value: V): Attributes<K,V> where K <: Hashable & Equatable<K>, V <: Equatable<V>
class GrpcError
/**
* grpc 错误处理类
*/
public open class GrpcError {
/*
* 空构造函数
*/
public init()
/*
* 构造函数
* 参数 code - 错误代码
* 参数 msg - 错误信息
*/
public init(code: UInt32, msg: String)
/*
* 判断是否是否有错误信息
* 返回值 Bool - 是否是空
*/
public func isNull(): Bool
/*
* 获取保存的错误信息
* 返回值 String - 错误信息
*/
public func message(): String
/*
* 获取格式化错误信息
* 返回值 String - 错误信息
*/
public open func getError(): String
/*
* 获取保存的错误代码
* 返回值 UInt32 - 代码
*/
public func getCode(): UInt32
}
class GrpcError
/**
* grpc 错误处理实现类
*/
public class ConnectionError <: GrpcError {
/*
* 空构造函数
*/
public init() {}
/*
* 构造函数
* 参数 code - 错误代码
*/
public init(code: UInt32)
/*
* 构造函数
* 参数 code - 错误信息
* 参数 temp - 缓存信息
* 参数 err - 错误对象
*/
public init(desc: String, temp: Bool, err: GrpcError)
/*
* 返回错误信息
* 参数 String - 错误信息
*/
public func error(): String
/*
* 返回错误缓存信息
* 参数 Bool - 错误缓存信息
*/
public func temporary(): Bool
/*
* 返回原始错误对象
* 参数 GrpcError - 返回错误对象
*/
public func origin(): GrpcError
/*
* 返回原始错误对象
* 参数 GrpcError - 返回错误对象
*/
public func unwrap(): GrpcError
}
class GrpcErrorMessage
/**
* grpc 错误处理实现类
*/
public class GrpcErrorMessage <: MessageLite {
/*
* 空构造函数
*/
public init() {}
/*
* 构造函数
* 参数 code - 错误代码
* 参数 msg - 错误信息
*/
public init(code: UInt32, msg: String)
/*
* 解码
* 参数 src - 源数据
*/
public func unpack<T>(src: T): Unit where T <: BytesReader
/*
* 编码
* 参数 out - 输出数据
*/
public func pack<T>(out: T): Unit where T <: BytesWriter
/*
* 转字符串
* 返回值 String - 转换的字符串
*/
public func toString(): String
/*
* 编码
* 参数 out - 输出数据
*/
public func packi(out: BytesWriter) { match (out) { case i_tmp: SimpleWriter => pack(i_tmp) case _ => pack(out) } }
/*
* 解码
* 参数 src - 源数据
*/
public func unpacki(src: BytesReader) { match (src) { case i_stmp: SimpleReader => unpack(i_stmp) case _ => unpack(src) } }
}
class Exponential
public class Exponential <: Strategy {
/*
* 空构造函数
*/
public init() {}
/*
* 构造函数
* 参数 config - 配置文件
*/
public init (config: Config)
/*
* 解码
* 参数 retries - 整数时间
* 参数 Duration - 返回时间
*/
public func Backoff(retries: Int64): Duration
}
class Exponential
public class Exponential <: Strategy {
/*
* 空构造函数
*/
public init() {}
/*
* 构造函数
* 参数 config - 配置文件
*/
public init (config: Config)
/*
* 解码
* 参数 retries - 整数时间
* 参数 Duration - 返回时间
*/
public func Backoff(retries: Int64): Duration
}
struct InsecureTC
public struct InsecureTC <: TransportCredentials {
/*
* 参数 s - 字符串
* 参数 socket - 一个 socket
* 返回值 (Socket, AuthInfo) - (Socket, AuthInfo)元组
*/
public func clientHandshake(s: String, socket: Socket): (Socket, AuthInfo)
/*
* 参数 socket - 一个 socket
* 返回值 (Socket, AuthInfo) - (Socket, AuthInfo)元组
*/
public func serverHandshake(socket: Socket): (Socket, AuthInfo)
/*
* 返回值 ProtocolInfo - 返回一个 ProtocolInfo
*/
public func info(): ProtocolInfo
/*
* 返回值 TransportCredentials - 返回一个 TransportCredentials
*/
public func clone(): TransportCredentials
/*
* 参数 s - 字符串
*/
public func overrideServerName(s: String):Unit
}
class MetaData
/*
* 初始化 MetaData
* 参数 m - 初始化的 HashMap
*/
public init(m: HashMap<String,ArrayList<String>>)
/*
* 按 key, value, key2, value2 顺序解析存入 MetaData, 该 api 为 MetaData 的核心 api
* 参数 kv - 要存入该 MetaData 的 ArrayList
* 返回值 mataData - 返回当前 MetaData
*/
public func pairs(kv: ArrayList<String>): MetaData
/*
* 返回该 MetaData 的大小
* 参数 kv - 要存入该 MetaData 的 ArrayList
* 返回值 Int64 - MetaData 的大小
*/
public func len(): Int64
/*
* 复制当前 MetaData
* 返回值 MetaData - 复制的新 MetaData 对象
*/
public func copy(): MetaData
/*
* 根据 key 取值
* 参数 k - key 值
* 返回值 ArrayList<String> - 返回对应的 value
*/
public func get(k: String): ArrayList<String>
/*
* 添加键值对
* 参数 k - key 值
* 参数 v - value 值
*/
public func set(k: String, v: ArrayList<String>): Unit
/*
* 添加键值对
* 参数 k - key 值
* 参数 v - value 值,已有值则叠加
*/
public func append(k: String, v: ArrayList<String>): Unit
/*
* 根据 key 值删除键值对
* 参数 k - key 值
*/
public func delete(k: String)
/*
* 两个 String 转一个元组
* 参数 k - key 值
* 参数 v - value 值
* 参数 (String, String) - 返回两个 String 的元组
*/
public func DecodeKeyValue(k: String, v: String): (String, String)
/*
* ArrayList 内多个 MetaData 合并存入一个 MetaData
* 参数 k - 要合并的 ArrayList<MetaData>
* 参数 MetaData - 合并完的 MetaData
*/
public func join(mds: ArrayList<MetaData>): MetaData
/*
* 定义 HashMap<String,ArrayList<String>> 类型为 MD
*/
type MD = HashMap<String,ArrayList<String>>
/*
* 从 Context 中读取 header 数值
* 参数 ctx - Context 值
* 返回值 (MD, Bool) - 返回读取到的值和是否成功
*/
public func fromIncomingContext(ctx: Context): (MD, Bool)
//Use the passed compressor or use the passed prepareMsg to return the hdr, payload, and data。
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func prepareMsg(m: MessageLite, codec: BaseCodec, flagCompress: Bool): (Array<Byte>, Array<Byte>, Array<Byte>, GrpcError)
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func encodeGrpc(c: BaseCodec, msg: MessageLite): (Array<Byte>, GrpcError)
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func compressGRPC(inbuf: Array<Byte>): (Array<Byte>, GrpcError)
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func msgHeader(buf: Array<Byte>, compData: Array<Byte>): (Array<Byte>, Array<Byte>)
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func recvAndDecompress(p: Parser, s: Stream, maxReceiveMessageSize: Int64, payInfo: PayloadInfo): (Array<UInt8>, GrpcError)
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func decompress(dcbuf: Array<UInt8>, maxReceiveMessageSize: Int64): (Array<Byte>, Int64, GrpcError)
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func recv(p: Parser, c: BaseCodec, s: Stream, m: MessageLite, maxReceiveMessageSize: Int64, payInfo: PayloadInfo): GrpcError
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 参数 flagCompress - 是否压缩
*/
public func toRPCErr(err: GrpcError): GrpcError
/*
* 准备数据
* 参数 k -数据
* 参数 codec - 编码
* 返回值 GrpcError - 错误
*/
public func defaultAnyAndServerStreamToGrpcError(any: Any, ss: ServerStream): GrpcError
class ServerWorkerData
public class ServerWorkerData{
/*
* 初始化
* 参数 st -数据
* 参数 stream - 编码
*/
public init(st: ServerTransport, stream: Stream)
/*
* 关闭
*/
public func close(): Unit
}
public func validate(md: HashMap<String, Array<String>>): String
class StreamDesc
public class StreamDesc {
/*
* 判空
* 参数 Bool -是否为空
*/
public func isNull(): Bool {
}
interface CallOption
public interface CallOption {
/*
* 执行前
* 参数 c -CallInfo
*/
func before(c: CallInfo): Unit
/*
* 执行后
* 参数 attempt -ClientStream
*/
func after(attempt: ClientStream): Unit
}
class BdpEstimator
public class BdpEstimator {
/*
* 判空
* 返回值 Bool -布尔值
*/
public func isNull(): Bool
/*
* 时间
* 参数 d - Array<UInt8> 值
*/
public func timesnap(d: Array<UInt8>): Unit
/*
* 添加
* 参数 n -UInt32值
*/
public func add(n: UInt32): Bool
}
class DefaultCbItem
public class DefaultCbItem <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class RegisterStream
public class RegisterStream <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class HeaderFrame
public class HeaderFrame <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class DataFrameCbItem
public class DataFrameCbItem <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class IncomingWindowUpdate
public class IncomingWindowUpdate <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class OutgoingWindowUpdate
public class OutgoingWindowUpdate <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class IncomingSettings
public class IncomingSettings <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class OutgoingSettings
public class OutgoingSettings <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class IncomingGoAway
public class IncomingGoAway <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class GoAway
public class GoAway <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class Ping
public class Ping <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class CleanupStream
public class CleanupStream <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
/*
* 返回值 Bool -布尔值
*/
public func isNull(): Bool {
}
class EarlyAbortStream
public class EarlyAbortStream <: CbItem {
/*
* 返回值 Bool -布尔值
*/
public func isTransportResponseFrame(): Bool
}
class ControlBuffer
public class ControlBuffer {
/*
* 返回值 Bool -布尔值
*/
public func finish(): Unit
}
class ItemList
public class ItemList<T> <: Collection<T> {
/*
* 清空 list
*/
public func clear(): Unit
/*
* 返回节点
* 返回值 ItemNode<T> - 返回节点
*/
public func dequeueAll(): ?ItemNode<T>
/*
* 判空
* 返回值 Bool - 布尔值
*/
public func isEmpty(): Bool
/*
* 返回迭代器
* 返回值 Iterator - 迭代器
*/
public func iterator(): Iterator<T>
/*
* 获取第一个值
* 返回值 T - 第一个值
*/
public func getFirst(): T
/*
* 获取第一个值
* 返回值 T - 第一个值
*/
public func peek(): T
/*
* 获取最后一个值
* 返回值 T - 最后一个值
*/
public func getLast(): T
/*
* 根据索引获取对应的值
* 参数 indexNum - 索引值
* 返回值 T - 获取的值
*/
public func get(indexNum: Int64): T
/*
* 根据索引设置值
* 参数 indexNum - 索引值
* 参数 element - 设置的值
* 返回值 T - 设置的值
*/
public func set(index: Int64, element: T): T
/*
* 预设值
* 参数 element - 设置的值
*/
public func prepend(element: T): Unit
/*
* 添加值
* 参数 element - 添加的值
*/
public func append(element: T): Unit
/*
* 排序
* 参数 element - 排序的值
*/
public func enqueue(element: T): Unit
/*
* 插入值
* 参数 indexNum - 插入的位置
* 参数 element - 设置的值
*/
public func insert(index: Int64, element: T): Unit
/*
* 移除值
* 参数 index - 移除的位置
* 返回值 T - 移除的值
*/
public func remove(index: Int64): T
/*
* 清空值
*/
public func removeFirst(): T
/*
* 队列化
*/
public func dequeue(): T
}
class LinkedListIterator
public class LinkedListIterator<T> <: Iterator<T> {
/*
* 获取下一个值
* 返回值 T - 下一个值
*/
public func next(): ?T
/*
* 获取迭代器
* 参数 Iterator<T> - 返回当前迭代器
*/
public func iterator(): Iterator<T>
}
/*
* 参数 flags -UInt8
* 返回值 Bool -布尔值
*/
public func hasPadded(flags: UInt8): Bool
/*
* 参数 flags -UInt8
* 返回值 Bool -布尔值
*/
public func hasPriority(flags: UInt8): Bool
/*
* 参数 flags -UInt8
* 返回值 Bool -布尔值
*/
public func isAck(flags: UInt8): Bool
/*
* 参数 flags -UInt8
* 返回值 Bool -布尔值
*/
public func headersEnded(flags: UInt8): Bool
/*
* 参数 flags -UInt8
* 返回值 Bool -布尔值
*/
public func streamEnded(flags: UInt8): Bool
/*
* 参数 hf -HeaderField
* 返回值 Bool -布尔值
*/
public func isPseudo(hf: HeaderField): Bool
public class GrpcDataFrame <: Frame {
/*
* 返回值 Bool -布尔值
*/
public func streamEnded(): Bool
}
class GrpcHeadersFrame
public class GrpcHeadersFrame <: Frame {
/*
* 参数 hf -HeaderField
* 返回值 Bool -布尔值
*/
public func setPriority(streamDependency: Int64, exclusive: Bool, weight: Int64)
/*
* 流结束
* 返回值 Bool -布尔值
*/
public func streamEnded(): Bool
/*
*流结束
* 返回值 Bool -布尔值
*/
public func headersEnded(): Bool
/*
* 存在私有值
* 返回值 Bool -布尔值
*/
public func hasPriority(): Bool
}
class GrpcSettingsFrame
public class GrpcSettingsFrame <: Frame{
/*
* 存在私有值
* 返回值 Bool -布尔值
*/
public func isAck(): Bool
/*
* 存在私有值
* 返回值 Bool -布尔值
*/
public func toString()
/*
* 存在私有值
* 返回值 Bool -布尔值
*/
public func setParameter(paramID: Int64, value: UInt32)
/*
* 存在私有值
* 返回值 UInt32 - UInt32
*/
public func getParameter(paramID: Int64): UInt32
}
class MetaHeadersFrame
public class MetaHeadersFrame <: Frame {
/*
* 判空
* 返回值 Bool -布尔值
*/
public func isNull(): Bool
/*
* 检查
* 返回值 GrpcError -GrpcError
*/
public func checkPseudos(): GrpcError
/*
* 变量检查
* 返回值 Array<HeaderField> -Array<HeaderField>值
*/
public func pseudoFields(): Array<HeaderField>
}
public func defaultStringToUnit(str: String): Unit
class Framer
public class Framer{
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeSettings(settings: HashMap<UInt16, UInt32>): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeSettingsAck(): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writePing(ack: Bool, data: Array<UInt8>): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeWindowUpdate(streamID: UInt32, incr: UInt32): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeHeaders(streamID: UInt32, chunk: Array<UInt8>, endStream: Bool,
endHeaders: Bool, padLength: Int64): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeRSTStream(streamID: UInt32, code: UInt32): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeContinuation(streamID: UInt32, endHeaders: Bool, headerBlockFragment: Array<UInt8>): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeGoAway(maxStreamID: UInt32, code: UInt32, debugData: Array<UInt8>): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeData(streamID: UInt32, endStream: Bool, data: Array<UInt8>): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func writeDataPadded(streamID: UInt32, endStream: Bool, data: Array<UInt8>, pad: Array<UInt8>): (Array<UInt8>, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func readFrame(): (Frame, Bool, GrpcError)
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func checkFrameOrder(frame: Frame): GrpcError
/*
* 写帧
* 参数 settings -HashMap<UInt16, UInt32>值
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError)值
*/
public func readMetaFrame(hf: GrpcHeadersFrame): (MetaHeadersFrame, Bool, GrpcError)
}
public func mapRecvMsgError(err: GrpcError): GrpcError
public func splitHostPort(strHostPort: String): (String, String, String)
class ServerHandlerTransport
public class ServerHandlerTransport <: ServerTransport{
/*
* 判空
* 返回值 Bool - Bool值
*/
public func isNull(): Bool
/*
* 创建新newServerHandlerTransport
* 返回值 (ServerTransport, GrpcError) - (ServerTransport, GrpcError) 值
*/
public static func newServerHandlerTransport(w: ResponseWriteStream, r: Request, stats: Array<GrpcHandler>): (ServerTransport, GrpcError)
/*
* 处理流
* 返回值 , GrpcError) - GrpcError 值
*/
public func handleStreams(startStream: (Stream) -> Unit, traceCtx: (String) -> Unit): Unit
/*
* 写header
* 返回值 , GrpcError) - GrpcError 值
*/
public func writeHeader(s: Stream, md: HashMap<String, Array<String>>): GrpcError
/*
* 写stream
* 返回值 , GrpcError) - GrpcError 值
*/
public func write(s: Stream, hdr: Array<UInt8>, data: Array<UInt8>): GrpcError
/*
* 写status
* 返回值 , GrpcError) - GrpcError 值
*/
public func writeStatus(s: Stream, st: GrpcError): GrpcError
/*
* 编码
* 参数 stream - stream值
*/
public func handleMarshal(err: GrpcError): (Array<UInt8>, GrpcError)
/*
* 写header
* 参数 stream - stream值
*/
public func writePendingHeaders(stream: Stream): Unit
/*
* 写commheader
* 参数 stream - stream值
*/
public func writeCommonHeaders(stream: Stream): Unit
/*
* 写customheader
* 参数 stream - stream值
*/
public func writeCustomHeaders(stream: Stream): Unit
/*
* 关闭
*/
public func close(): Unit
/*
* 获取值
* 返回值 SocketAddress - SocketAddress值
*/
public func getRemoteAddr(): SocketAddress
/*
* 获取值
*/
public func incrMsgSent(): Unit
/*
* 获取值
*/
public func incrMsgRecv(): Unit
/*
* 获取值
*/
public func drain(): Unit
/*
* 返回当前的 Context
* 返回值 Context - 当前的 Context
*/
public func context(): Context
}
class GRPCConn
public enum GRPCConn {
/*
* 获取值
*/
public func read(buf: Array<Byte>): Int64
/*
* 获取值
*/
public func read(buf: Array<Byte>, timeout: Duration): Int64
/*
* 写值
*/
public func write(buf: Array<Byte>): Unit
/*
* 写值
*/
public func write(buf: Array<Byte>, timeout: Duration): Unit
/*
* 关闭
*/
public func close(): Unit
/*
* 获取地址
*/
public func getLocalAddr(): SocketAddress
/*
* 获取地址
*/
public func getRemoteAddr(): SocketAddress
}
class GRPCServer
public enum GRPCServer{
/*
* 获取地址
* 返回值 GRPCConn - GRPCConn 值
*/
public func accept(): GRPCConn
/*
* 获取地址
* 返回值 GRPCConn - GRPCConn 值
*/
public func accept(timeout: Duration): GRPCConn
/*
* 获取地址
* 返回值 String - String 值
*/
public func toString(): String
/*
* 获取地址
* 返回值 SocketAddress - SocketAddress 值
*/
public func getAddr(): SocketAddress
/*
* 关闭
*/
public func close(): Unit
/*
* 判断是否关闭
* 返回值 Bool - Bool 值
*/
public func isClosed(): Bool
}
class BufWriter
public class BufWriter {
/*
* 写值
* 返回值 Bool - Bool 值
* 返回值 (Int64, GrpcError) - (Int64, GrpcError) 值
*/
public func write(src: Array<UInt8>): (Int64, GrpcError)
/*
* 刷新流
* 返回值 GrpcError - GrpcError 值
*/
public func flush(): GrpcError
/*
* 读取值
* 返回值 (Int64, GrpcError) - (Int64, GrpcError) 值
*/
public func read(srcRead: Array<UInt8>): (Int64, GrpcError)
}
/*
* 超时时间设置
* 返回值 (Duration, Bool) - (Duration, Bool) 值
*/
public func timeoutUnitToDuration(char: Char): (Duration, Bool)
/*
* 超时时间设置
* 返回值 (Duration, GrpcError) - (Duration, GrpcError) 值
*/
public func decodeTimeout(s: String): (Duration, GrpcError)
/*
* 头数据设置
* 返回值 Bool - Bool) 值
*/
public func isWhitelistedHeader(hdr: String): Bool
/*
* 头编码
* 返回值 (String, GrpcError) - (String, GrpcError) 值
*/
public func decodeMetadataHeader(k: String, v: String): (String, GrpcError)
/*
* 二进制编码
* 返回值 (Array<UInt8>, GrpcError) - (Array<UInt8>, GrpcError) 值
*/
public func decodeBinHeader(str: String): (Array<UInt8>, GrpcError)
class GRPCConnSocket
public class GRPCConnSocket {
/*
* 设置超时时间
* 返回值 SocketAddress -SocketAddress 值
*/
public func setDeadline(time: Duration): Unit
/*
* socket地址获取
* 参数 SocketAddress -SocketAddress 值
*/
public func read(buffer: Array<Byte>): Int64
/*
* 写值
* 参数 buffer -Array<Byte> 值
*/
public func write(buffer: Array<Byte>): Unit
/*
* 关闭
*/
public func close(): Unit
/*
* 获取远程地址
* 返回值 SocketAddress -SocketAddress 值
*/
public func getLocalAddr(): SocketAddress
/*
* 获取远程地址
* 返回值 SocketAddress -SocketAddress 值
*/
public func getRemoteAddr(): SocketAddress
}
class GRPCServerSocket
public class GRPCServerSocket <: Hashable & Equatable<GRPCServerSocket> {
/*
* 获取哈希值
* 返回值 Int64 -Int64 值
*/
public func hashCode(): Int64
/*
* 获取远程地址
* 参数 time -Duration 值
*/
public func setDeadline(time: Duration): Unit
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func accept(): GRPCConn
/*
* 获取远程地址
* 返回值 SocketAddress -SocketAddress 值
*/
public func getAddr(): SocketAddress
/*
* 判断是否关闭
*/
public func isClosed(): Unit
/*
* 关闭
*/
public func close(): Unit
}
class Http2Server
public class Http2Server <: ServerTransport {
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func outgoingGoAwayHandler(g: GoAway): (Bool, GrpcError)
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func handleStreams(handle: (Stream) -> Unit, traceCtx: (String) -> Unit): Unit
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func writeHeader(s: Stream, md: HashMap<String, Array<String>>): GrpcError
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func write(s: Stream, hdr: Array<UInt8>, data: Array<UInt8>): GrpcError
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func writeStatus(s: Stream, st: GrpcError): GrpcError
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func checkForHeaderListSize(it: Any): Bool
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func close(): Unit
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func getRemoteAddr(): SocketAddress
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func drain(): Unit
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func incrMsgSent(): Unit
/*
* 获取远程地址
* 返回值 GRPCConn -GRPCConn 值
*/
public func incrMsgRecv(): Unit
/*
* 返回当前的 Context
* 返回值 Context - 当前的 Context
*/
public func context(): Context
}
class LoopyWriter
public class LoopyWriter {
public func handle(any: CbItem): GrpcError
public func run(): GrpcError
}
class OutStream
public class OutStream{
/*
* 判空
* 返回值 Bool -Bool 值
*/
public func isNull(): Bool
/*
* 删除自身
*/
public func deleteSelf(): Unit
}
class OutStreamList
public class OutStreamList {
public func enqueue(s: OutStream): Unit
public func dequeue(): OutStream
}
class ServerInHandle
public class ServerInHandle {
public func isNull(): Bool
}
public interface RPCStats {
func isRPCStats(): Unit
func isClient(): Bool
}
class RecvBuffer
public class RecvBuffer {
public func put(rMsg: RecvMsg): Unit
}
class RecvBufferReader
public class RecvBufferReader <: InputStream {
public func read(buffer: Array<Byte>): Int64
}
class TransportReader
public class TransportReader <: InputStream {
public func read(buffer: Array<Byte>): Int64
}
class Stream
public open class Stream <: InputStream {
/*
* 转CommonAuthInfo
* 参数 Bool - Bool
*/
public func isNull(): Bool
/*
* 转CommonAuthInfo
* 参数 Int64 - Int64
*/
public func read(data: Array<Byte>): Int64
/*
* 转CommonAuthInfo
* 参数 data - Array<Byte>
*/
public func write(data: Array<Byte>): Unit
/*
* 转CommonAuthInfo
* 参数 rMsg - RecvMsg
*/
public func write(rMsg: RecvMsg): Unit
/*
* 获取getTrailer
* 参数 CommonAuthInfo - CommonAuthInfo
*/
public func getTrailer(): HashMap<String, Array<String>>
/*
* 获取ContentSubtype
* 参数 String - String
*/
public func getContentSubtype(): String
/*
* 获取方法
* 参数 String - String
*/
public func getMethod(): String
/*
* 获取状态
* 返回值 GrpcError - GrpcError
*/
public func getStatus(): GrpcError
/*
* 获取返回值
* 返回值 String - String
*/
public func getRecvCompress(): String
/*
* 转CommonAuthInfo
* 参数 st - UInt32
*/
public func swapState(st: UInt32): UInt32
/*
* 转CommonAuthInfo
* 参数 md - HashMap<String
*/
public func headerAddMD(md: HashMap<String, Array<String>>): Unit
}
class ServerTransport
public abstract class ServerTransport <: Hashable & Equatable<ServerTransport>{
public func handleStreams(handle: (Stream) -> Unit, traceCtx: (String) -> Unit): Unit
public func writeHeader(s: Stream, md: HashMap<String, Array<String>>): GrpcError
public func write(s: Stream, hdr: Array<UInt8>, data: Array<UInt8>): GrpcError
public func writeStatus(s: Stream, st: GrpcError): GrpcError
public func close(): Unit
public func getRemoteAddr(): SocketAddress
public func drain(): Unit
public func incrMsgSent(): Unit
public func incrMsgRecv(): Unit
public func hashCode(): Int64
}
public struct CommonAuthInfo {
/*
* 转CommonAuthInfo
* 参数 CommonAuthInfo - CommonAuthInfo
*/
public func getCommonAuthInfo(): CommonAuthInfo
}
public enum SecurityLevel {
/*
* 转字符串
* 参数 String - String
*/
public func toString(): String
}
public class Info <: AuthInfo {
/*
* any转unit
* 参数 any - any
*/
public func authType(): String
}
public class BaseCodec {
/*
* 编码
* 参数 m - 数据
*/
public func marshal(m: MessageLite): Array<Byte>
/*
* 解码
* 参数 data - 数据
*/
public func unmarshal(data: ArrayList<Byte>, m: MessageLite): Unit
}
/*
* any转unit
* 参数 any - any
*/
public func defaultAnyToUnit(any: Any): Unit {}
/*
* error转any
* 参数 any - any
*/
public func defaultErrorToUnit(any: Any): Unit
/*
* uint32转unit
* 参数 num - UInt32
*/
public func defaultUInt32ToUnit(num: UInt32): Unit
/*
* UInt64转unit
* 参数 num - UInt64
*/
public func defaultUInt64ToUnit(num: UInt64): Unit
/*
* 解压
* 参数 num - Int64
*/
public func defaultInt64ToUnit(num: Int64): Unit
/*
* 解压
* 参数 it - Any
*/
public func defaultAnyToBool(it: Any): Bool
/*
* unit转unit
*/
public func defaultUnitToUnit(): Unit
/*
* any转grpcerr
* 参数 Any - 原始值
*/
public func defaultAnyToGrpcError(it: Any): GrpcError
/*
* unit32转error
* 参数 it - 原始值
* 返回值 GrpcError - 错误值
*/
public func defaultUInt32ToGrpcError(it: UInt32): GrpcError
/*
* bytearra转unit
* 参数 src - 原始值
*/
public func defaultByteArrayStreamToUnit(buf: ByteArrayStream): Unit
/*
* string转error
* 参数 src - 原始值
*/
public func defaultServerInHandle(info: String): GrpcError
class DEFAULT_MESSAGELITE
public class DEFAULT_MESSAGELITE <: MessageLite{
/*
* 解压
* 参数 src - 原始值
*/
public func unpacki(src: BytesReader)
/*
* 压缩
* 参数 out - 解码值
*/
public func packi(out: BytesWriter)
/*
* 转字符串
* 返回值 String - 字符串
*/
public func toString(): String
}
public func defaultAnyToAnyAndGrpcError(any: Any): (MessageLite, GrpcError)
public func contentType(contentSubtype: String): String
class EmptyResponseWriteStream
public class EmptyResponseWriteStream <: ResponseWriteStream {
/*
* 获取header值
* 返回值 Header - 获取值
*/
public func header(): Header
/*
* 写值
* 参数 buf - 原始值
*/
public func write(buf: Array<UInt8>): Unit
/*
* 写状态码
* 参数 statusCode - 原始值
*/
public func writeStatusCode(statusCode: Int64): Unit
/*
* 刷新流
*/
public func flush(): Unit
}
/*
* 转换
* 参数 data - 原始值
* 返回值 UInt32 - 转换值
*/
public func getUInt32(data: Array<UInt8>): UInt32
/*
* 转换
* 参数 data - 原始值
*/
public func putUInt16(buf: Array<Byte>, num: UInt16, index: Int64): Unit
/*
* 转换
* 参数 data - 原始值
*/
public func putByte3(buf: Array<Byte>, num: UInt32, index: Int64): Unit
/*
* 转换
* 参数 data - 原始值
*/
public func putUInt32(buf: Array<Byte>, num: UInt32, index: Int64): Unit
/*
* 转换
* 参数 data - 原始值
*/
public func putByteArray(buf: Array<Byte>, src: Array<Byte>, index: Int64): Unit
/*
* 转换
* 参数 data - 原始值
* 返回值 Bool - 是否成功
*/
public func isReservedHeader(hdr: String): Bool
/*
* 开启选项
* 返回值 Bool - 是否开启
*/
public func isOn(): Bool
/*
* 编码值
* 参数 arr - 原始值
* 返回值 String - 编码值
*/
public func parseUInt8ToHexadecimal(num: UInt8): String
/*
* 编码值
* 参数 arr - 原始值
* 返回值 String - 编码值
*/
public func encodeGrpcMessage(msg: String): String
/*
* 编码值
* 参数 arr - 原始值
* 返回值 String - 编码值
*/
public func encodeGrpcMessageUnchecked(str: String): String
/*
* 编码值
* 参数 arr - 原始值
* 返回值 String - 编码值
*/
public func join(mds: HashMap<String, Array<String>>, md: HashMap<String, Array<String>>): Unit
/*
* 编码值
* 参数 arr - 原始值
* 返回值 String - 编码值
*/
public func appendHeaderFieldsFromMD(headerFields: Array<HeaderField>, md: HashMap<String, Array<String>>): Array<HeaderField>
/*
* 编码值
* 参数 arr - 原始值
* 返回值 String - 编码值
*/
public func encodeMetadataHeader(k: String, v: String): String
/*
* 编码值
* 参数 arr - 原始值
* 返回值 String - 编码值
*/
public func encodeBinHeader(arr: Array<UInt8>): String
/*
* 返回纳秒
* 返回值 Int64 - 当前纳秒
*/
public func unixNano(): Int64
/*
* streamid 检查
* 参数 streamID - streamid 值
* 返回值 Bool - 是否正常
*/
public func validStreamID(streamID: UInt32): Bool
/*
* 头数据检查
* 参数 v - 头数据
* 返回值 Bool - 是否正常
*/
public func validWireHeaderFieldName(v: String): Bool
/*
* contenttype检查
* 参数 contentType - 数据
* 返回值 (String, Bool) - 返回元组
*/
public func handleContentSubtype(contentType: String): (String, Bool)
/*
* 清空heafer数据
* 参数 contentType - 数据
* 返回值 String - 清空
*/
public func clearHeadAndTailControlCharacters(contentType: String): String
/*
* 检查时候能打印
* 参数 msg - 数据
* 返回值 Bool - 是否能打印
*/
public func hasNotPrintable(msg: String): Bool
/*
* 添加 hashmap 值
* 参数 str - 数据
*/
public func appendHashMapValue(map: HashMap<String, Array<String>>, name: String, value: String): Unit
/*
* 数据编码
* 参数 str - 数据
* 返回值 (Int32, Int64) - 编码后的值
*/
public func decodeRuneInString(str: String): (Int32, Int64)
class Proto
class Proto {
/*
* 数据编码
* 参数 m - 数据
* 返回值 Array<Byte> - 编码后的值
*/
public static func marshal(m: MessageLite): Array<Byte>
/*
* 数据解码
* 参数 b - 加密数据
* 参数 m - 解码后的值
*/
public static func unmarshal(b: ArrayList<Byte>, m: MessageLite): Unit
}
class BlockByteQueue
public class BlockByteQueue {
/*
* 初始化队列
* 参数 size -初始化大小
*/
public init(size: Int64)
/*
* 添加所有值
* 参数 buf -要添加的值
*/
public func appendAll(buf: Array<UInt8>): Unit
/*
* 获取值
* 参数 buf - 获取的 Array<UInt8>
* 返回值 Int64 - 获取的大小
*/
public func get(buf: Array<UInt8>): Int64
/*
* 获取所有值
* 返回值 Array<UInt8> - 获取到的值
*/
public func getAll(): Array<UInt8>
/*
* 清空所有值
*/
public func clear()
}
class WriteQuota
public class WriteQuota {
/*
* 判空
* 返回值 Bool - 是否为空
*/
public func isNull(): Bool {
/*
* 获取值
* 参数 sz - 索引值
* 返回值 GrpcError - 错误
*/
public func get(sz: Int32): GrpcError {
/*
* 补充值
* 参数 n - Int64 值
*/
public func realReplenish(n: Int64): Unit {
}
class TrInFlow
内部 dns 解析器
class TrInFlow {
/*
* 添加新行
* 参数 n - UInt32 值
* 返回值 UInt32 -添加的值
*/
public func newLimit(n: UInt32): UInt32
/*
* 重置 TrInFlow
* 参数 UInt32 - 重置数量
*/
public func reset(): UInt32
/*
* 更新windowssize
*/
public func updateEffectiveWindowSize(): Unit
}
class dnsBuilder
内部 dns 解析器
public open class dnsBuilder <: resolver.Builder & resolver.Resolver {
/**
* 初始化状态
* 参数 scheme - 解析头字符串
*/
public func initialState(s: resolver.State)
/**
* 传入解析的目标,构造解析器
* 参数 target - 解析目标
* 返回值 resolver.Resolver - 返回对应的解析器
* 返回值 GrpcError - 返回错误信息
*/
public func build(target: resolver.Target): (resolver.Resolver, GrpcError)
/**
* 返回当前解析器的头
* 参数 String - 返回解析头字符串
*/
public func scheme(): String
/**
* 用户传入自定义参数做自定义解析
* 参数 o - ResolveNowOptions 类型的自定义参数
*/
public func resolveNow(o: resolver.ResolveNowOptions): Unit
/**
* 调用用户自定义的 ClientConn 的 close
*/
public func close(): Unit
/**
* 调用用户自定义的 ClientConn 的 updateState
* 参数 resolver.State - 要处理的 State
* 返回值 GrpcError - 返回一个 GrpcError
*/
public func updateState(s: resolver.State): GrpcError
/**
* 处理 grpc 错误,调用用户自定义的 ClientConn 的 reportError
* 参数 err - 要处理的 GrpcError
*/
public func reportError(err: GrpcError): Unit
}
interface dnsBuilder
Context上下文传值
public interface Context {
/**
* 根据 key 获取值
* 参数 key - 键
* 返回值 Any - 值
*/
func value(key: Any): Any
}
public class emptyCtx <: Context{
/**
* 根据 key 获取值,表示空Context的emptyCtx的该函数默认抛NoneValueException异常
*/
public func value(_: Any): Any
}
public class valueCtx <: Context{
public var context: Context
/**
* 初始化 valueCtx
* 参数 context - 上级 Context
* 参数 key - 键
* 参数 val - 值
*/
public init(context: Context, key: Any, val: Any)
/**
* 根据 key 获取值
* 参数 key - 键
* 返回值 Any - 值
*/
public func value(key: Any): Any
}