rap(par[::-1]) 是高级且快速的python异步rpc
项目描述
说唱
rap(par[::-1]) 是高级且快速的python异步rpc
rapmsgpack通过和Python asyncio多路conn实现非常快速的通信,同时支持高并发。通过 Python 函数和 TypeHint实现protobufof 。Grpc
注意:当前rapAPI 在后续版本中可能会发生重大变化
说唱首版功能思路来自aiorpc
1.安装
pip install rap
2.快速入门
服务器
import asyncio
from typing import AsyncIterator
from rap.server import Server
def sync_sum(a: int, b: int) -> int:
return a + b
async def async_sum(a: int, b: int) -> int:
await asyncio.sleep(1) # mock io
return a + b
async def async_gen(a: int) -> AsyncIterator[int]:
for i in range(a):
yield i
loop = asyncio.new_event_loop()
rpc_server = Server() # init service
# register func
rpc_server.register(sync_sum)
rpc_server.register(async_sum)
rpc_server.register(async_gen)
# run server
loop.run_until_complete(rpc_server.create_server())
try:
loop.run_forever()
except KeyboardInterrupt:
# stop server
loop.run_until_complete(rpc_server.await_closed())
客户
客户端支持通过raw_callandcall方法调用服务,但是这样不能充分使用TypeHint的功能,建议使用@client.register注册函数后再调用。
注意:和rap.client没有区别,但是注册的函数用户可以直接使用,所以修饰的函数应该类似于:async defdef@client.register@client.register
async def demo(): pass
例子:
import asyncio
from typing import AsyncIterator
from rap.client import Client
client: "Client" = Client() # init client
# Declare a function with no function. The function name, function type and return type must be the same as the server side function (async def does not differ from def)
def sync_sum(a: int, b: int) -> int:
pass
# The decorated function must be an async def function
@client.register()
async def sync_sum(a: int, b: int) -> int:
pass
# The decorated function must be the async def function, because the function is a generator syntax, to `yield` instead of `pass`
@client.register()
async def async_gen(a: int) -> AsyncIterator:
yield
async def main():
await client.connect()
# Call the call method; read the function name and then call `raw_call`.
print(f"call result: {await client.call(sync_sum, 1, 2)}")
# Basic calls to rap.client
print(f"raw call result: {await client.raw_call('sync_sum', 1, 2)}")
# Functions registered through `@client.register` can be used directly
# await async_sum(1,3) == await client.raw_call('async_sum', 1, 2)
# It is recommended to use the @client.register method, which can be used by tools such as IDE to determine whether the parameter type is wrong
print(f"decorator result: {await sync_sum(1, 3)}")
async_gen_result: list = []
# Example of an asynchronous generator, which by default opens or reuses the current session of the rap (about the session will be mentioned below)
async for i in async_gen(10):
async_gen_result.append(i)
print(f"async gen result:{async_gen_result}")
asyncio.run(main())
三、功能介绍
3.1.注册功能
服务端支持defand async def,如果是def函数,多线程运行。注册时会检查函数的参数和返回值的TypeHints,如果类型与json指定的类型不匹配会报错。
服务器带有一个注册库。如果同组有重复注册,会报错。可以使用group参数定义要注册的组,也可以用参数重新定义注册的名称name(客户端调用时还需要指定对应的组)。
另外is_private注册的时候可以设置为True,这样该函数只能被本地的rap.client调用。
import asyncio
from typing import AsyncIterator
from rap.server import Server
def demo1(a: int, b: int) -> int:
return a + b
async def demo2(a: int, b: int) -> int:
await asyncio.sleep(1)
return a + b
async def demo_gen(a: int) -> AsyncIterator[int]:
for i in range(a):
yield i
server: Server = Server()
server.register(demo1) # register def func
server.register(demo2) # register async def func
server.register(demo_gen) # register async iterator func
server.register(demo2, name='demo2-alias') # Register with the value of `name`
server.register(demo2, group='new-group') # Register and set the groups to be registered
server.register(demo2, group='root', is_private=True) # Register and set the group to be registered, and set it to private
对于客户端,建议使用,client.register代替。
使用 Python 语法定义函数名称、参数、参数类型和返回值类型,它允许调用者像调用普通函数一样调用函数,并且可以使用 TypeHint 功能通过工具检查函数。注意:使用时,一定要使用。client.callclient.raw_callclient.registerclient.registerasync def ...
from typing import AsyncIterator
from rap.client import Client
client: Client = Client()
# register func
@client.register()
async def demo1(a: int, b: int) -> int: pass
# register async iterator fun, replace `pass` with `yield`
# Since `async for` will make multiple requests to the same conn over time, it will check if the session is enabled and automatically reuse the current session if it is enabled, otherwise it will create a new session and use it.
@client.register()
async def demo_gen(a: int) -> AsyncIterator: yield
# Register the general function and set the name to demo2-alias
@client.register(name='demo2-alias')
async def demo2(a: int, b: int) -> int: pass
# Register the general function and set the group to new-group
@client.register(group='new-group')
async def demo2(a: int, b: int) -> int: pass
3.2.会话
rap客户端支持会话功能,开启会话后,所有请求都只会通过当前会话的conn请求到对应的服务器,而每一个请求,header中的session_id都会设置当前会话id,方便服务器识别。
rap会话支持显式和隐式设置,每种设置都有自己的优点和缺点,没有强制性限制。
from typing import AsyncIterator
from rap.client import Client
client = Client()
def sync_sum(a: int, b: int) -> int:
pass
@client.register()
async def async_sum(a: int, b: int) -> int:
pass
@client.register()
async def async_gen(a: int) -> AsyncIterator[int]:
yield
async def no_param_run():
# The rap internal implementation uses the session implicitly via the `contextvar` module
print(f"sync result: {await client.call(sync_sum, 1, 2)}")
print(f"async result: {await async_sum(1, 3)}")
# The asynchronous generator detects if a session is enabled, and if so, it automatically reuses the current session, otherwise it creates a session
async for i in async_gen(10):
print(f"async gen result:{i}")
async def param_run(session: "Session"):
# By explicitly passing the session parameters in
print(f"sync result: {await client.call(sync_sum, 1, 2, session=session)}")
print(f"sync result: {await client.raw_call('sync_sum', 1, 2, session=session)}")
# May be a bit unfriendly
print(f"async result: {await async_sum(1, 3, session=session)}")
# The asynchronous generator detects if a session is enabled, and if so, it automatically reuses the current session, otherwise it creates a session
async for i in async_gen(10):
print(f"async gen result:{i}")
async def execute(session: "Session"):
# The best way to call a session explicitly, using a method similar to the mysql cursor
# execute will automatically recognize the type of call
print(f"sync result: {await session.execute(sync_sum, arg_list=[1, 2])}")
print(f"sync result: {await session.execute('sync_sum', arg_list=[1, 2])}")
print(f"async result: {await session.execute(async_sum(1, 3))}")
# The asynchronous generator detects if a session is enabled, and if so, it automatically reuses the current session, otherwise it creates a session
async for i in async_gen(10):
print(f"async gen result:{i}")
async def run_once():
await client.connect()
# init session
async with client.session as s:
await no_param_run()
await param_run(s)
await execute(s)
await client.await_close()
3.3.渠道
channel支持双工方式的client-server交互,类似于Http的WebSocket,需要注意的是channel不支持group设置。
客户端仅@client.register支持注册通道功能,其特点是类型为单个参数Channel。通道将保持一个会话,并且只会在通道启用到关闭之间通过 conn 与服务器通信。为了避免使用 'while True',通道支持使用 'async for' 语法和使用 'while await channel.loop()' 语法而不是 'while True
from rap.client import Channel, Client
from rap.client.model import Response
client = Client()
@client.register()
async def async_channel(channel: Channel) -> None:
await channel.write("hello") # send data
cnt: int = 0
while await channel.loop(cnt < 3):
cnt += 1
print(await channel.read_body()) # read data
@client.register()
async def echo_body(channel: Channel) -> None:
await channel.write("hi!")
# Reads data, returns only when data is read, and exits the loop if it receives a signal to close the channel
async for body in channel.iter_body():
print(f"body:{body}")
await channel.write(body)
@client.register()
async def echo_response(channel: Channel) -> None:
await channel.write("hi!")
# Read the response data (including header data), and return only if the data is read, or exit the loop if a signal is received to close the channel
async for response in channel.iter_response():
response: Response = response # help IDE check type....
print(f"response: {response}")
await channel.write(response.body)
3.4.ssl
Python asyncio由于模块的高度封装,rap可以很方便的和ssl一起使用
# Quickly generate ssl.crt and ssl.key
openssl req -newkey rsa:2048 -nodes -keyout rap_ssl.key -x509 -days 365 -out rap_ssl.crt
客户端.py
from rap.client import Client
client = Client(ssl_crt_path="./rap_ssl.crt")
服务器.py
from rap.server import Server
rpc_server = Server(
ssl_crt_path="./rap_ssl.crt",
ssl_key_path="./rap_ssl.key",
)
3.5.事件
服务器端分别支持启动前start_event和stop_event关闭后的事件处理。
from rap.server import Server
async def mock_start():
print('start event')
async def mock_stop():
print('stop event')
# example 1
server = Server(start_event_list=[mock_start()], stop_event_list=[mock_stop()])
# example 2
server = Server()
server.load_start_event([mock_start()])
server.load_stop_event([mock_stop()])
3.6.中间件
rap目前支持2种中间件::
- conn中间件:在创建conn时使用,比如限制链接总数等...参考block.py,该
dispatch方法会传入一个conn对象,然后根据规则判断是否释放(返回await self .call_next(conn)) 或拒绝它(等待 conn.close) - 消息中间件:只支持普通函数调用(不支持
Channel),类似使用starlette中间件参考access.py 消息中间件会传入4个参数:request(当前请求对象), call_id(当前调用id), func(当前调用function)、param(当前参数)和request返回call_id和result(函数执行结果或异常对象)
另外,中间件支持start_event_handle和stop_event_handle方法,分别在Server启动和关闭时调用。
例子:
from rap.server import Server
from rap.server.middleware import AccessMsgMiddleware, ConnLimitMiddleware
rpc_server = Server()
rpc_server.load_middleware([ConnLimitMiddleware(), AccessMsgMiddleware()])
3.7.处理器
rap处理器用于处理入站和出站流量,其中用于process_request入站流量和process_response出站流量。
rap.client和处理器的方法rap.server基本相同,rap.server支持start_event_handle和stop_event_handle方法,分别在Server启动和关闭时调用
客户端负载处理器示例
from rap.client import Client
from rap.client.processor import CryptoProcessor
client = Client()
client.load_processor([CryptoProcessor('key_id', 'xxxxxxxxxxxxxxxx')])
服务器负载处理器示例
from rap.server import Server
from rap.server.processor import CryptoProcessor
server = Server()
server.load_processor([CryptoProcessor({'key_id': 'xxxxxxxxxxxxxxxx'})])
4.插件
middlewarerap通过and支持插件功能processor,middleware只支持服务端,processor支持客户端和服务端
4.1.加密传输
加密传输只对请求和响应的正文内容进行加密,对header等内容不进行加密。加密时,增加nonce参数防止重放,增加timestamp参数防止超时访问。
客户端示例:
from rap.client import Client
from rap.client.processor import CryptoProcessor
client = Client()
# The first parameter is the id of the secret key, the server determines which secret key is used for the current request by the id of the secret key
# The second parameter is the key of the secret key, currently only support the length of 16 bits of the secret key
# timeout: Requests that exceed the timeout value compared to the current timestamp will be discarded
# interval: Clear the nonce interval, the shorter the interval, the more frequent the execution, the greater the useless work, the longer the interval, the more likely to occupy memory, the recommended value is the timeout value is 2 times
client.load_processor([CryptoProcessor("demo_id", "xxxxxxxxxxxxxxxx", timeout=60, interval=120)])
服务器示例:
from rap.server import Server
from rap.server.processor import CryptoProcessor
server = Server()
# The first parameter is the secret key key-value pair, key is the secret key id, value is the secret key
# timeout: Requests that exceed the timeout value compared to the current timestamp will be discarded
# nonce_timeout: The expiration time of nonce, the recommended setting is greater than timeout
server.load_processor([CryptoProcessor({"demo_id": "xxxxxxxxxxxxxxxx"}, timeout=60, nonce_timeout=120)])
4.2. 限制最大conn数
仅限服务器端使用,可以限制服务器端的最大链接数,超过设置值将不处理新请求
from rap.server import Server
from rap.server.middleware import ConnLimitMiddleware, IpMaxConnMiddleware
server = Server()
server.load_middleware(
[
# max_conn: Current maximum number of conn
# block_timeout: Access ban time after exceeding the maximum number of conn
ConnLimitMiddleware(max_conn=100, block_time=60),
# ip_max_conn: Maximum number of conn per ip
# timeout: The maximum statistics time for each ip, after the time no new requests come in, the relevant statistics will be cleared
IpMaxConnMiddleware(ip_max_conn=10, timeout=60),
]
)
4.3.限制ip访问
支持限制单个ip或整段ip,同时支持白名单和黑名单模式,如果启用白名单,则默认禁用黑名单模式
from rap.server import Server
from rap.server.middleware import IpBlockMiddleware
server = Server()
# allow_ip_list: whitelist, support network segment ip, if filled with allow_ip_list, black_ip_list will be invalid
# black_ip_list: blacklist, support network segment ip
server.load_middleware([IpBlockMiddleware(allow_ip_list=['192.168.0.0/31'], block_ip_list=['192.168.0.2'])])
5.高级功能
TODO , 这个功能还没有实现
6.协议设计
TODO , 正在编辑文档
7.说唱运输介绍
TODO , 正在编辑文档