跳转到内容
快速开始

基于 Mirobody 进行开发

开发一个 Provider

写一个 BasePullProvider 子类丢进 providers/:工厂方法、ProviderInfo 元数据、凭据校验、OAuth 授权流程、定时拉取,以及归一到 StandardPulseData。

给 Mirobody 接入一个健康数据源,只需要写一个 BasePullProvider 子类,放进引擎启动时会扫描的目录里就行。没有编译步骤,不用改注册表,也不用碰任何核心文件:加载器会找到你的类、调用它的工厂方法,然后把返回值注册进去。

一份能运行起来的源码检出

一个装好 pip install -e . 的虚拟环境,加上 docker compose up -d pg redis 起来的 Postgres 和 Redis。如果你的数据源走 OAuth,Redis 不是可选项:临时授权状态就存在那里。见开发环境搭建

数据源的 API 文档

基础 URL、认证方式、你需要的数据对应的端点,以及响应结构。你写的每一行都是对这些载荷的一次变换,因此应尽早保存几条真实响应:它们就是你后面的测试夹具。

一个放凭据的地方

client id、secret 和基础 URL 都通过 safe_read_cfg("YOUR_KEY") 从配置里读,不要直接读环境变量。键名里含 _KEY_PASSWORD_PASS_PWD_SECRET_SK_TOKEN 的值会被自动加密落盘。见配置

启动时,引擎遍历 PROVIDER_DIRS 的每一项,它默认是 mirobody/pulse/providers 加上仓库根目录的 providers/。把你自己的 provider 放进 providers/,升级就永远不会碰到它:

有四条规则决定你的 provider 到底存不存在:加载器先用 mirobody_*/provider_*.py 做 glob,再去找一个继承 BasePullProvider 的类(类名只用来预筛,要求以 Provider 结尾):

规则取值
目录mirobody_<slug>/
模块provider_<something>.py
<Something>Provider(BasePullProvider)
工厂create_provider(config) 返回实例,或返回 None 表示保持关闭

代码树里最小的一个完整 provider 是 PostgreSQL 那个:不调任何厂商 API,只实现了必需的那一套接口,可以当模板读。

必须由你实现的只有三个:infosave_raw_data_to_dbis_data_already_processed。其余的要么在基类上有可用的默认实现,要么在缺失时抛一个把你漏写的方法名点出来的异常。

create_provider(config) 是个 classmethod,它只做一件事:判断这个 provider 在当前环境下能否工作。 返回 None 是官方支持的「不进注册表」的方式:一个没配好的集成就这样消失掉,而不是等到请求时才报错。

两种写法分别适用于不同情况:集成本身可选时,用开关键把住;集成必须有凭据才能工作时,用凭据本身把住。ENABLE_PGSQL_DEVICE 并不在随仓库发布的 config.yaml 里:不加上这个键,PostgreSQL provider 就会一直是关着的。

mirobody_acme/provider_acme.py
class AcmeProvider(BasePullProvider):
@classmethod
def create_provider(cls, config: dict) -> Optional["AcmeProvider"]:
client_id = config.get("ACME_CLIENT_ID")
if not client_id: # 没配置 → 这个 provider 不存在
return None
return cls(config)

info 返回一个 ProviderInfo,并且不许碰网络:列出所有 provider 本来就该是零成本的。真正值得琢磨的两个字段是 auth_typeconnect_info_fields,它们是绑在一起的。

LinkType.CUSTOMIZED,你描述一个表单、前端把它渲染出来;填进去的值会以 credentials["connect_info"] 的形式回到你手上。PostgreSQL provider 声明了五个字段(用户名、密码、host、端口、数据库),这里只摘了两个。

一个 OAuth provider 一个字段都不声明:它没有表单,只有一次跳转。

PASSWORDCUSTOMIZED 两类 provider,BasePullProvider.link() 会在写入任何数据之前调用 _validate_credentials_v2(credentials),而抛异常就是全部的失败协议。「你填的不对」抛 ValueError,「数据源连不上」抛 RuntimeError,消息会一路回到调用方。

OAuth provider 走的是另一条路:link() 返回一个给浏览器打开的 link_web_url,连接由后面的回调路由收尾。两个 OAuth 版本在代码树里都有实例(Whoop 和 Oura 是 OAuth 2.0,Garmin 是 OAuth 1.0a),而且两者都不把状态放在进程内存里。握手期的临时状态存在 Redis 里,所以浏览器可以回到一个跟发起流程时不同的实例上。

别自己手搓这套流程:引擎自带一个可复用的 OAuth 2.0 客户端,provider 用组合而不是继承的方式接它。

有些数据源在刷新授权上偏离 RFC 6749,refresh_extra_params 就是留给它们的口子:Whoop 要求刷新时重新带上 scope,那就放进这里,而不是去 fork 一份客户端。

第一阶段拼出授权 URL,并把握手状态写进 Redis,TTL 取自 OAUTH_TEMP_TTL_SECONDS(默认 900 秒)。

state 里带着调用方的 return_url,好让浏览器事后能被送回原处;redirect_uri 跟它存在一起,是因为换令牌时必须把授权时用过的同一个值交给令牌端点。

第二阶段换取令牌。exchange_code_for_tokens 会取走并清掉那份握手状态,以 authorization_code 授权类型换令牌,算出过期时间并加密落库。你的 provider 只负责接上两头。

刷新令牌也不是你的活。需要令牌时就调 get_valid_access_token(user_id, provider_slug, db_service):还剩五分钟以上寿命就直接返回已存的那个,否则花掉 refresh token、把新的一对存下来、再把新的 access token 交给你。如果一个 refresh token 都没存,它返回 None,那就是「用户必须重新授权」的信号。

OAuth 1.0a 多一次往返、而且每个请求都要签名,所以它不复用上面那个客户端。引擎里有一份可照抄的实现:Garmin Provider 示例。你要提供的仍然只有一个 callback 方法,参数是 (oauth_token, oauth_verifier)

你不需要加路由。一条路由服务所有 provider,它按 info.auth_type 分派。

所以你的义务是提供一个参数个数对得上的 callback 方法:OAuth 2.0 是 (code, state),OAuth 1.0a 是 (oauth_token, oauth_verifier)。这条路由是 GET /api/v1/pulse/{platform}/{provider}/callback。调用方没给 return_url 时,它返回一小段把弹窗关掉、并通知打开方的 HTML。连接流程的完整用法见使用 Provider

端点 URL 都带默认值,所以真正必填的只有三个 secret。随仓库发布的模板里,Garmin 与 Whoop 的授权 URL、令牌 URL 与 API base URL 都已填好,空着的只有:

config.yaml
GARMIN_CLIENT_ID: ""
GARMIN_CLIENT_SECRET: ""
GARMIN_REDIRECT_URL: ""
WHOOP_CLIENT_ID: ""
WHOOP_CLIENT_SECRET: ""
WHOOP_REDIRECT_URL: ""

给你自己的数据源沿用同一套形状:<SLUG>_CLIENT_ID<SLUG>_CLIENT_SECRET<SLUG>_REDIRECT_URL,以及可覆盖的 <SLUG>_AUTH_URL / <SLUG>_TOKEN_URL / <SLUG>_API_BASE_URL。因为键名以 _SECRET 结尾,client secret 会被自动加密落盘。Oura 也遵守这个约定,尽管它的三个键并不在模板里:读一个不存在的键只会获得空值,工厂方法于是拒绝构造这个 provider。

对推送式的数据源,同样没有什么要注册的。两条路由接收载荷,都会分派到你的 provider:

Terminal window
# explicit provider — preferred
POST /api/v1/pulse/providers/theta_garmin/webhook
# universal: provider is sniffed out of the body's `source` field
POST /api/v1/pulse/providers/webhook

msg_id 取自请求头 Svix-Id,头缺失时回退成一个格式化的时间戳。你的 is_data_already_processed 和你的存储表就是靠它认出一次重复投递的。

如果你的数据源只能被轮询,那就明说,并实现一个方法。register_pull_task 返回 True(基类默认值)就会给你一个带分布式锁的定时任务;返回 False 就是退出,Garmin 和 PostgreSQL 两个 provider 都这么做:Garmin 是因为它的 webhook 已经完成了这项工作,PostgreSQL 是因为它只校验一个连接。

循环由基类替你驱动:它载入每个已连接用户的凭据,按用户调你的 pull_from_vendor_api,跳过被 is_data_already_processed 否掉的,剩下的推进写入路径。你要返回的是一个自带说明的包裹列表,一种数据一个,像 Whoop provider 那样。

那个 data_type 就是你后面 format_data_v2 里用来分派的标签,所以名字定一次、两处都用同一套。凭据不是「用户名 + 密码」的数据源有两个挂钩点:用你自己的签名覆写 pull_from_vendor_api,再覆写 _pull_and_push_for_user(credentials) 去调它:Whoop 就是这样传 access token 和 refresh token 而不是用户名密码的。按 slug 配置的节奏和锁时长见 Pulse Provider 体系;没有单独配置的 slug 是每小时一拉、锁 30 分钟。

这是只有你能写的部分,也是整个 provider 契约里唯一的硬约束:不管你的数据源说什么形状的话,format_data_v2 必须返回一个 StandardPulseData。下游的一切(聚合、指标检索、agent 的健康工具)只读这个模型,不读别的。

输入是一个 FormatDataInput,它被刻意拆成两半:一个基类已经从数据库里解析好的 context(内部用户 id、厂商侧用户 id、时区、消息 id),和原封不动的 payload。因为身份和时区都是预先解析好的,format_data_v2 完全不做 I/O,是一个纯函数,因此可以离线测试。

每条记录的 type 必须是一个已登记的指标名,不是你数据源的字段名:Oura 戒指和 Garmin 手表的读数正是靠这一点才能放在一起比较。所以真正的工作是一张从厂商字段到 StandardIndicator 的表,把它当数据维护,不要散落在大量 if 分支中。Oura 的写法最简:源单位已经跟标准单位一致时就写一个裸的指标,需要写入路径去换算时就写一个 (indicator, source_unit) 元组。

"total_sleep_duration"
(StandardIndicator.DAILY_TOTAL_SLEEP_TIME, "s")
"efficiency"
StandardIndicator.SLEEP_EFFICIENCY
"average_hrv"
StandardIndicator.HRV_RMSSD
"steps"
StandardIndicator.DAILY_STEPS
"spo2_percentage.average"
StandardIndicator.BLOOD_OXYGEN
Oura 的映射表:左边是厂商字段,右边是已登记的指标(单位不同时附上源单位)。

Whoop 的表是同一个思路加上一个显式的换算器,因为它的载荷说的是毫秒和千焦:每一项是 (indicator_name, converter, unit)。指标名从 registry 里挑(见健康指标);如果登记表里没有合适的,去扩登记表,别在这里现编一个 type 字符串。

mirobody_acme/provider_acme.py
def format_data_v2(self, raw: dict, user_id: str) -> StandardPulseData:
records = [
StandardPulseRecord(
source="acme",
type="heartrate", # 已登记的指标名,不是 Acme 自己的字段名
value=sample["bpm"],
unit="/min",
timestamp=sample["epoch_ms"],
)
for sample in raw["samples"]
]
return StandardPulseData(userId=user_id, healthData=records)

还剩两个抽象方法,它们管的是持久化而不是变换。save_raw_data_to_db 把载荷原样写下来(于是一个格式化 bug 只意味着重运行一次,而不是数据丢了),并返回载荷里每个用户一条,一个批量 webhook 就是这样扇出成好几次 format_data_v2 调用的。is_data_already_processed 是推送前的幂等闸门。

列名不是随便起的。get_table_name() 默认是 health_data_<slug 去掉 "theta_">get_user_id_column() 默认是 theta_user_idget_query_columns() 默认是 id、用户 id 列、external_user_idmsg_idraw_datacreate_atupdate_atis_del。照这个布局来,管理控制台就能免费地翻页浏览并重新格式化你存下的载荷;不照,就把那三个方法一起覆写掉。

没有构建步骤:重启进程、读日志。每一次成功装载都会输出一行日志;而 GET /api/v1/pulse/providers 不需要令牌,所以它是确认你的类被找到了的最快方式:

/path/to/providers
mirobody serve
# → Loaded provider from /path/to/providers/mirobody_acme/provider_acme.py
# → ✅ Loaded provider: [theta_acme]
curl -s http://localhost:18080/api/v1/pulse/providers

如果 slug 不在这个响应里,原因必是四者之一,而日志会告诉你是哪一个:目录名或模块名没匹配上 glob、类没有继承 BasePullProvider、导入抛了异常,或者 create_provider 因为缺一个配置键而返回了 None

然后把变换本身认真验一遍(夹具、快照,以及线上路由),见 Provider 测试

随包发布的四个 provider 可以照着读:mirobody/pulse/providers/