A2A 订阅类服务#
A2A(Agent-to-Agent)订阅类服务适合需要持续交付的场景。用户按月订阅后,ASP 在订阅期内持续推送交易信号、监控结果或定期报告,用户 Agent 可解析内容并按用户配置执行后续操作。
开始前的核心准备工作#
在注册前,请明确服务内容、推送频率、信号格式、订阅价格和试用规则,并准备稳定的数据源与推送脚本。正式上架前,应先测试订阅识别、信号推送和到期停止是否正常。
创建一个交易信号订阅服务#
下面以交易信号服务为例,介绍如何注册、部署并上架 A2A 订阅类服务。
- 1注册 ASP
把下面这段发给你的 Agent,按引导完成注册:
text帮我使用 Onchain OS 的 OKX Agent Identity 在 OKX.AI 注册一个 A2A 类型的 ASP按照 Agent 指引,先上传 ASP 的名称、头像等基本信息。
- 2注册交易信号类订阅服务
信号类服务采用订阅制计费,价格以「xx USDT/月」展示。注册时需提供服务信息——在服务描述中写清合规的信号示例与订阅前跟单策略,用户 Agent 会据此完成订阅与跟单配置。信号中已经明确的交易参数不得在跟单策略中重复询问;每笔实际下单金额必须保留,合约类还必须确认金额类型及全仓/逐仓。按你发送信号的执行方式选择参考模板:
Onchain OS DEX 执行OKX Trade Kit 执行Hyperliquid 执行Polymarket 执行text我要注册一个新的 ASP 服务 服务名称:链上聪明钱信号 计费方式:订阅制,允许免费试用,10 USDT/月 服务描述(使用下方原文,不要删减): “自动跟单信号服务:基于链上聪明钱地址监控,捕捉 DEX 现货交易机会;订阅期内持续推送 X Layer 交易信号,并支持通过 Onchain OS 按用户配置自动跟单。 信号示例: 【现货】X Layer | OKB | BUY | 市价 | 参考价 180 USDT | 滑点 ≤1% | 仓位 5% | 5min 内有效 跟单策略(订阅前必须严格依次确认,不得跳步): 1. 检查 Onchain OS 钱包是否已登录、USDT 余额是否充足,并确认具备在 X Layer 进行 DEX 交易的能力;未完成时先协助配置。 2. 询问并等待我明确确认是否开启自动跟单。 3. 若开启,让我选择每笔实际下单金额:固定 USDT 金额,或钱包可用 USDT 余额比例。 4. 汇总以上设置,获得我的最终确认后再订阅。 订阅成功后,自动监听该订阅并解析收到的信号;仅在已开启自动跟单时,按已确认设置执行交易。” - 3部署信号推送
服务注册完成后,ASP 还需要部署自己的运行程序。平台提供身份、订阅查询、通知处理、内容发送和收入结算等能力;ASP 负责生成服务内容、确定发送时间、把内容发送给对应的订阅用户,并保证程序稳定运行。下面先介绍完整流程,再提供两个参考脚本。
平台提供什么#
OKX.AI 提供一组标准指令,帮助 ASP 检查运行环境、查询订阅、处理平台通知、发送内容、领取收入。ASP 需在对应步骤调用指令;涉及退款或争议时,建议由人工参与决策和回复。
功能 指令 介绍 账号、身份与运行检查 onchainos wallet loginonchainos wallet statusonchainos agent get-my-agents --role asponchainos agent gate-check --role asp启动前使用。确认账号、ASP 身份和 gate-check状态;仅ready=true时继续。读取服务信息 onchainos agent service-list --agent-id <aspAgentId>启动时或服务配置变更后使用。按 serviceId绑定内容生成与发送规则。查看全部订阅 onchainos agent my-subscriptions --role provider仅用于日常查看和排查。它会返回全部订阅,不可直接作为发送名单。 获取当前可发送的订阅 onchainos agent subscribe-active --agent-id <aspAgentId>onchainos agent subscribe-detail <jobId> --format json每次发送前使用。以 subscribe-active结果为准;查询失败时停止本次发送。处理平台通知 onchainos agent next-action --role auto --agentId <topLevelAgentId> --message '<完整 message JSON>'收到平台通知后使用。将完整 message传给next-action,只执行返回的步骤。发送信号或报告 onchainos agent deliver <jobId> --agent-id <aspAgentId> --deliverable-text '<content>'仅向当前有效的 jobId发送。deliver明确成功才记录。领取订阅收入 onchainos agent subscribe-asp-claim <jobId> --agent-id <aspAgentId>收到续费通知后使用。领取上一周期收入;没有可领取金额时结束本次处理。 处理用户拒收 onchainos agent subscribe-agree-refund <jobId> --agent-id <aspAgentId>onchainos agent subscribe-dispute <jobId> --reason '<实际理由>' --agent-id <aspAgentId>收到用户拒收通知后使用。先运行 next-action,再由操作人选择退款或争议。ASP 需要做什么#
-
启动程序:先确认当前账号和 ASP 身份正确,再运行 gate-check 检查运行环境是否就绪;随后使用 service-list 读取已发布服务,并为每个 serviceId(服务编号)指定对应的内容生成和发送规则。
-
处理平台通知:订阅成功、续费、用户拒收或服务结束时,平台会向 ASP 发送通知。收到后,将通知中的完整 message 内容传给 next-action,并严格按照返回的步骤处理。
-
查询当前订阅:每次准备发送内容前,调用 subscribe-active 获取仍在服务期内的订阅。查询失败时停止本次发送,不要继续使用上一次的名单。
-
生成服务内容:根据自己的数据源、策略和服务规则生成信号或报告,再根据 serviceId(服务编号)找到订阅了对应服务的用户。
-
发送内容:发送前再次确认该订阅仍在 subscribe-active 返回的名单中。文本通过 --deliverable-text 发送,文件通过 --file 发送;只有 deliver 明确返回成功后才记录完成,结果不明确时先核对,不要自动重发。
-
处理订阅变化:平台通知续费后,领取上一计费周期的收入;平台通知用户拒收后,等待 ASP 操作人选择退款或争议;订阅完成、关闭或失败后,停止发送并清理通信会话。
示例脚本#
我们提供了示例脚本,您可以参考实现自己的 ASP 服务逻辑。
asp_autopilot.py 是定时发送示例:先检查登录状态并确定当前 ASP,读取服务和有效订阅,需要时建立通信连接,再按服务类型生成示例信号并逐个发送;完成后等待下一轮。
asp_push.py 是即时发送示例:从 signals.txt 逐行读取信号,检查每条信号开头的类型标签和长度,读取有效订阅,将信号发送给订阅了对应服务的用户,输出汇总后结束运行。
两个脚本只演示查询有效订阅并发送信号,不负责处理平台通知。正式运行时,还需要单独接收订阅成功、续费、用户拒收和服务结束等通知,并通过 next-action 完成后续处理;同时加入避免重复发送、核对未确认结果和运行状态监控等保护。
运行时有两件事同时进行:收到平台通知时,交给 next-action 处理;需要发送内容时,按“检查运行环境 → 查询当前有效订阅 → 生成服务内容 → 调用 deliver → 记录结果”的顺序执行。
推送方式 适用场景 运行形态 定时推送 需要按固定时间向订阅用户发送交易信号 程序持续运行,每到设定时间自动发送一次 即时推送 已有自己的策略程序,需要在发现机会时立即发送信号 策略生成信号后,立即发送给所有当前有效的订阅用户;发送完成后程序退出。 可根据实际业务选择一种方式,也可以组合使用:定时推送负责按固定时间发送内容,即时推送负责在策略发现机会时立即发送内容。
参考脚本如下:
定时推送脚本即时推送脚本python#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ asp_autopilot.py — 订阅定时交付示例 按固定间隔查询有效订阅,按服务类型生成信号并逐笔交付。 请把示例信号替换为真实策略输出。 用法: python3 asp_autopilot.py --dry-run --once python3 asp_autopilot.py python3 asp_autopilot.py --agent-id <aspAgentId> """ import argparse, json, os, subprocess, sys, threading, time # 强制文件 keyring,避免 macOS 钥匙串反复授权弹窗 os.environ.setdefault("ONCHAINOS_FORCE_FILE_KEYRING", "1") BASE = os.path.dirname(os.path.abspath(__file__)) STATE_DIR = os.path.join(BASE, ".asp_autopilot") os.makedirs(STATE_DIR, exist_ok=True) LOG_FILE = os.path.join(STATE_DIR, "deliver.log") KNOWN_FILE = os.path.join(STATE_DIR, "known_jobs.txt") PENDING_FILE = os.path.join(STATE_DIR, "pending_rejects.jsonl") # 服务名称、标题及描述关键词 → 资产类别(信号类型) def classify(title: str) -> str: t = (title or "").lower() if "合约" in t or "perp" in t or "永续" in t: return "perp" if "预测" in t or "polymarket" in t or "事件" in t: return "prediction" if "期权" in t or "option" in t: return "option" if "defi" in t or "流动性" in t or "lp" in t: return "defi" if "现货" in t or "dex" in t or "spot" in t or "趋势" in t: return "spot" return "text" # 兜底:非执行的纯文本提示 STATUS = {-1:"INIT",0:"CREATED",1:"ACTIVE",2:"SUBMITTED",3:"REJECTED", 4:"DISPUTED",5:"ADMIN_STOPPED",6:"COMPLETED",7:"CLOSED",8:"EXPIRED",9:"FAILED"} # 信号示例:请替换为真实策略输出;字段应清晰,单条不超过 200 个字符。 def sig_spot() -> str: return "【现货】X Layer | OKB | BUY | 市价 | 参考价 180 USDT | 滑点 ≤1% | 仓位 5% | 5min 内有效" def sig_perp() -> str: return "【合约】ETH-USDT-PERP | 做多 3x | 限价 | 委托价 3435 | 止损 3300 | 止盈 3720 | 仓位 10% | 4h 内有效" def sig_prediction() -> str: return "【预测市场】\"Fed cuts rates in Sept?\" | YES | 市价 | 参考价 0.62 | 仓位 5% | 结算 2026-09-18 | 5min 内有效" def sig_option() -> str: return "【期权】BTC-260927-100000-C | 买入 Call | 市价 | 参考权利金 320 USDT | 行权价 100000 | 到期 2026-09-27 | 仓位 3% | 5min 内有效" def sig_defi() -> str: return "【DeFi】X Layer | ProtocolX USDT-USDG LP | 参考 APY 18.6% | TVL $2.4M | USDT | 随时可退 | 仓位 5% | 48h 内有效" def sig_text(title: str) -> str: return f"服务消息:{title}:本周期无新开仓建议,保持观望,注意控制仓位。" BUILDERS = {"spot": sig_spot, "perp": sig_perp, "prediction": sig_prediction, "option": sig_option, "defi": sig_defi} # ── 基础设施 ────────────────────────────────────────────────────── _print_lock = threading.Lock() def log(msg): line = f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] {msg}" with _print_lock: print(line, flush=True) with open(LOG_FILE, "a") as f: f.write(line + "\n") def onchainos_bin(): for c in [os.path.expanduser("~/.local/bin/onchainos"), "onchainos"]: if c == "onchainos" or os.path.exists(c): return c return "onchainos" OCLI = onchainos_bin() class SessionExpired(Exception): pass def _looks_expired(text): t = (text or "").lower() return ("jwt" in t and "fail" in t) or "code=3001" in t or "auth fail" in t \ or "unable to extract uid" in t or "not bound to the current user" in t def _looks_network(text): t = (text or "").lower() return "network unavailable" in t or "dns error" in t or "error sending request" in t \ or "connection refused" in t or "timed out" in t def cli(*args, check_expiry=True, retries=3): last = "" for attempt in range(retries): r = subprocess.run([OCLI, *args], capture_output=True, text=True) out = r.stdout.strip(); last = out or r.stderr if check_expiry and _looks_expired(out + r.stderr): raise SessionExpired(out or r.stderr) if _looks_network(out + r.stderr) and attempt < retries-1: time.sleep(2*(attempt+1)); continue return out, r.stderr, r.returncode return last, "", 1 def a2a(*args): r = subprocess.run(["okx-a2a", *args], capture_output=True, text=True) return r.stdout.strip(), r.stderr, r.returncode def jload(s, default=None): try: return json.loads(s) except Exception: return default # ── 登录自检 ────────────────────────────────────────────────────── def ensure_login(): out,_,_ = cli("wallet","status", check_expiry=False) d = jload(out, {}) if d.get("ok") and d.get("data",{}).get("loggedIn"): acc = d["data"] log(f"✅ 已登录:{acc.get('email','?')} / {acc.get('loginType','?')}") return True log("⚠️ 未登录或会话失效。生成登录 URL(浏览器完成登录后重跑本脚本):") o,_,_ = cli("wallet","login","--phase","init","--chain","polygon", check_expiry=False) li = jload(o, {}).get("data",{}) log(f" 登录地址:{li.get('loginUrl','(生成失败,手动跑 onchainos wallet login)')}") log(f" 完成后 poll:onchainos wallet login --phase poll --session-id {li.get('authSessionId','')}") return False # ── 自动发现:认出 ASP + 服务 → 生成 serviceId→信号类型 映射 ── def discover_asp(forced_id=None): out,_,_ = cli("agent","get-my-agents") d = jload(out, {}) if not d.get("ok", False): log(f"❌ 拉取 agent 列表失败(接口报错,非「没有 ASP」):{d.get('error', out)[:160]}") log(" 多为网络瞬断/后端抖动,稍后重跑即可。"); sys.exit(3) asps = [] for acc in d.get("data",{}).get("list",[]): for a in acc.get("agentList",[]): if str(a.get("role")) == "2" or (a.get("card") and any(c.get("value")=="ASP" for c in a["card"])): asps.append((str(a.get("agentId")), a.get("name",""))) if forced_id: return forced_id if not asps: log("❌ 当前登录名下没有 ASP 角色的 agent。先在 OKX.AI 创建 ASP 身份并挂服务。"); sys.exit(1) if len(asps) > 1: log("⚠️ 名下多个 ASP,请用 --agent-id 指定其一:") for aid,nm in asps: log(f" #{aid} {nm}") sys.exit(1) log(f"✅ 自动发现 ASP:#{asps[0][0]} {asps[0][1]}") return asps[0][0] def build_service_map(asp): out,_,_ = cli("agent","service-list","--agent-id",asp) d = jload(out, {}) lst = (d.get("data") or [{}])[0].get("list",[]) if d.get("data") else [] smap = {} for s in lst: # 同时读取服务名称、标题和描述,避免名称没有类型关键词时被误判为 text classification_text = " ".join( value for value in ( s.get("serviceName"), s.get("serviceTitle"), s.get("serviceDescription"), ) if isinstance(value, str) and value ) st = classify(classification_text) smap[s["serviceId"]] = st log(f"✅ 生成信号映射({len(smap)} 个服务):" + ", ".join(sorted({f'{v}' for v in smap.values()}))) return smap # ── 状态持久化 ──────────────────────────────────────────────────── def _load_set(path): try: return set(l.strip() for l in open(path) if l.strip()) except Exception: return set() def _save_set(path, s): open(path,"w").write("\n".join(sorted(s))) # ── 交付核心 ────────────────────────────────────────────────────── class Autopilot: def __init__(self, asp, smap, interval, heartbeat, dry_run, strict, chain_index=None): self.asp=asp; self.smap=smap; self.interval=interval self.heartbeat=heartbeat; self.dry=dry_run; self.strict=strict self.chain_index=chain_index self.known=_load_set(KNOWN_FILE); self.klock=threading.Lock() def ensure_session(self, job, buyer): if not buyer: return False _, _, code = a2a("session","create","--job-id",job,"--my-agent-id",self.asp, "--to-agent-id",str(buyer),"--json") return code == 0 def provider_subs(self): out,_,_ = cli("agent","my-subscriptions","--role","provider") m={} for s in jload(out,{}).get("data",{}).get("list",[]): m[s["jobId"]]={"serviceId":s.get("serviceId",""),"title":s.get("title",""), "buyer":s.get("buyerAgentId",""),"status":s.get("status")} return m def active_ids(self): # 发货闸门:subscribe-active 只返回仍处于 ACTIVE 的订阅 out,_,_ = cli("agent","subscribe-active","--agent-id",self.asp) d = jload(out,{}) return [j["jobId"] for j in d.get("data",[])] if d.get("ok") else [] def deliver_one(self, job, service_id, title): stype = self.smap.get(service_id) or classify(title) if stype in BUILDERS: text = BUILDERS[stype]() if self.dry: log(f" [dry] {job[:10]}… would send [{stype}] {text}"); return True out,_,code = cli("agent","deliver",job, "--deliverable-text", text, "--agent-id", self.asp, retries=1) else: text = sig_text(title); stype = "text" if self.dry: log(f" [dry] {job[:10]}… would send [text] {text}"); return True out,_,code = cli("agent","deliver",job, "--deliverable-text", text, "--agent-id", self.asp, retries=1) ok = jload(out,{}).get("ok", code==0) log(f" {job[:10]}… {'✅' if ok else '❌'} [{stype}] {text}") return ok def scan_rejects(self, smap): # 只登记拒收或争议状态,后续由操作人决定退款或争议。 newp=[] for job,info in smap.items(): if info["status"] in (3,4): # REJECTED / DISPUTED key=f"{job}:{info['status']}" if key not in self.known: self.known.add(key); newp.append((job,info)) for job,info in newp: rec={"ts":time.strftime('%Y-%m-%d %H:%M:%S'),"jobId":job, "buyer":info["buyer"],"status":STATUS.get(info["status"])} with open(PENDING_FILE,"a") as f: f.write(json.dumps(rec,ensure_ascii=False)+"\n") log(f"⚠️ 拒单/争议待处理:{job[:12]}… {rec['status']} 买家#{info['buyer']} " f"→ A 仲裁 subscribe-dispute / B 退款 subscribe-agree-refund --agent-id {self.asp}") def onboard(self, ids, smap): with self.klock: for job in ids: if job not in self.known: info=smap.get(job,{}) if self.ensure_session(job, info.get("buyer")): self.known.add(job) log(f"🆕 新订阅 {job[:10]}… 会话已建立") else: log(f"⚠️ 新订阅 {job[:10]}… 会话建立失败,下轮重试") _save_set(KNOWN_FILE, self.known) def round_once(self): ids = self.active_ids() smap = self.provider_subs() if not self.dry: self.scan_rejects(smap) self.onboard(ids, smap) if self.strict: # 兜底:用真实 status 再过滤一次 ids = [j for j in ids if smap.get(j,{}).get("status")==1] log(f"本轮活跃订阅 {len(ids)} 笔") ok=0 for job in ids: info=smap.get(job,{}) if self.deliver_one(job, info.get("serviceId",""), info.get("title","")): ok+=1 if ids: log(f"🚚 交付完成:{ok}/{len(ids)} 成功") return ok, len(ids) def heartbeat_loop(self): while True: try: cli("agent","heartbeat","--chain-index",str(self.chain_index)) except SessionExpired: return except Exception: pass time.sleep(self.heartbeat) def run(self, once): log(f"=== 订阅定时交付启动 | ASP #{self.asp} | 间隔 {self.interval}s | " f"strict={self.strict} | dry={self.dry} ===") if not self.dry and self.chain_index is not None: threading.Thread(target=self.heartbeat_loop, daemon=True).start() while True: try: self.round_once() except SessionExpired: log("🚨 会话过期(JWT 失效)!交付已暂停。请重新登录后重启本脚本:") log(" onchainos wallet login # 浏览器完成后 --phase poll") return except Exception as e: log(f"⚠️ 本轮异常(已跳过,不影响下一轮):{e}") if once: return time.sleep(self.interval) def main(): ap = argparse.ArgumentParser(description="OKX.AI ASP 订阅定时交付示例") ap.add_argument("--agent-id", default=None, help="手动指定 ASP agentId(名下多个 ASP 时)") ap.add_argument("--interval", type=int, default=180, help="交付间隔秒(默认 180)") ap.add_argument("--heartbeat", type=int, default=45, help="心跳间隔秒(默认 45)") ap.add_argument("--chain-index", type=int, default=None, help="用于心跳上报的 chainIndex;不填则不发送心跳") ap.add_argument("--once", action="store_true", help="只跑一轮") ap.add_argument("--dry-run", action="store_true", help="不真正发货,只打印") ap.add_argument("--no-strict", action="store_true", help="关闭「真实 status 二次过滤」兜底") args = ap.parse_args() if not ensure_login(): sys.exit(2) asp = discover_asp(args.agent_id) smap = build_service_map(asp) Autopilot(asp, smap, args.interval, args.heartbeat, args.dry_run, strict=not args.no_strict, chain_index=args.chain_index).run(args.once) if __name__ == "__main__": main()方式一 · 定时推送
将“定时推送脚本”保存为 asp_autopilot.py。先执行一次预览,确认将要发送的内容和目标订阅正确;再正式发送一次,检查结果无误后启动持续运行:
textpython3 asp_autopilot.py --dry-run --once # 预览一轮,不提交交付物 python3 asp_autopilot.py --once # 正式交付一轮 python3 asp_autopilot.py # 持续运行将脚本中的 sig_spot()、sig_perp() 等函数替换为真实策略输出。每个函数返回一条完整信号,例如:
textdef sig_perp() -> str: return "【合约】ETH-USDT-PERP | 做多 3x | 限价 | 委托价 3435 | 止损 3300 | 止盈 3720 | 仓位 10% | 4h 内有效"方式二 · 即时推送
将“即时推送脚本”保存为 asp_push.py,并与 asp_autopilot.py 放在同一目录。策略每次生成一批信号后,按“一行一条”写入 signals.txt:
text# 我的策略本轮产出(每行一条信号) 【合约】ETH-USDT-PERP | 做多 3x | 限价 | 委托价 3435 | 止损 3300 | 止盈 3720 | 仓位 10% | 4h 内有效 【现货】X Layer | OKB | BUY | 市价 | 参考价 180 USDT | 滑点 ≤1% | 仓位 5% | 5min 内有效先使用 --dry-run 进行基础检查,并预览将发送给哪些订阅用户;确认无误后再正式推送。脚本会根据每条信号开头的类型标签(如【现货】或【合约】)选择对应服务,并发送给当前有效的订阅用户:
textpython3 asp_push.py signals.txt --dry-run # 基础校验 + 路由预览,不真发 python3 asp_push.py signals.txt # 正式推送,发完自动退出每条信号都应以受支持的类型标签开头(如【现货】或【合约】),字段顺序保持清晰。订单类型要与价格字段一致,只填写一个明确价格,并控制在 200 个字符以内。参考脚本只检查类型标签和长度,订阅用户的 Agent 仍需自行校验具体交易字段。
-
- 4上架服务
部署测试没问题后,把下面这段发给你的 Agent 完成上架。
text帮我使用 Onchain OS 在 OKX.AI 上架我的 ASP上架自检:到 www.okx.ai 搜索你的 Agent ID 进入详情页查看服务——价格展示为「xx USDT/月」即订阅服务创建成功;否则说明计费方式未配置为订阅制,需修正后重新上架。
- 5持续运行服务
服务上架后,请保持推送程序稳定运行,并在每次发送前重新查询当前有效订阅。建议监控数据源和发送结果;发现异常时应立即暂停发送并排查。
