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、新订阅自动建会话、发送信号、心跳保活全部自动完成。两种推送方式的区别只在「信号什么时候发出去」:
推送方式 适用场景 运行形态 定时推送 固定周期信号,定期给订阅用户推送交易信号 脚本常驻运行,每隔固定时间自动发一轮信号 即时推送 已有策略系统,信号由策略触发 策略产出信号时批量推给所有活跃订阅者,发完即退 两者也可按需搭配使用:平时定时推送兜底,策略出信号时即时推送。
完整参考脚本如下:
定时推送脚本即时推送脚本python#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ asp_autopilot.py — OKX.AI ASP 交易信号持续交付「一站式」守护进程(V2·信号类型版) ================================================================================ 一条命令搞定:登录自检 → 自动发现 ASP 身份与服务 → 自动生成「服务→信号类型」映射 → 监听 + 心跳保活 + 持续交付(v1.2 规范信号)+ 拒单登记。 python3 asp_autopilot.py # 自动发现 ASP + 启动持续交付 python3 asp_autopilot.py --once # 只跑一轮(冒烟测试) python3 asp_autopilot.py --dry-run # 不真正发货,只打印本轮将交付什么 python3 asp_autopilot.py --interval 60 # 交付间隔秒数(默认 180) python3 asp_autopilot.py --agent-id 4941 # 手动指定 ASP(名下多个 ASP 时) 依赖:全局 onchainos CLI + okx-a2a CLI。Python3 标准库即可。 """ 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") SEQ_FILE = os.path.join(STATE_DIR, "seq.txt") 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"} # ══════════════════════════════════════════════════════════════════ # 信号文本模板 —— 换成自己的策略时,必须遵循交易信号 v1.2: # (a) 使用合法头部与固定字段顺序 (b) 订单类型与价格字段匹配,且只写一个具体价格 # (c) 仓位使用 N% (d) 单条信号不超过 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 _read_int(path, d=0): try: return int(open(path).read().strip()) except Exception: return d 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))) SEQ_LOCK = threading.Lock() def next_delivery_id(): with SEQ_LOCK: n = _read_int(SEQ_FILE) + 1 open(SEQ_FILE,"w").write(str(n)) return f"{time.strftime('%Y%m%d')}-{n:05d}" # ── 交付核心 ────────────────────────────────────────────────────── class Autopilot: def __init__(self, asp, smap, interval, heartbeat, dry_run, strict): self.asp=asp; self.smap=smap; self.interval=interval self.heartbeat=heartbeat; self.dry=dry_run; self.strict=strict self.known=_load_set(KNOWN_FILE); self.klock=threading.Lock() # 外部买家必须先建 XMTP 会话,否则 deliver 上链成功但 P2P 推送失败 def ensure_session(self, job, buyer): if not buyer: return a2a("session","create","--job-id",job,"--my-agent-id",self.asp, "--to-agent-id",str(buyer),"--json") 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) did = next_delivery_id() 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, ) 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) ok = jload(out,{}).get("ok", code==0) log(f" {job[:10]}… {'✅' if ok else '❌'} [{stype}] {text}") return ok def scan_rejects(self, smap): # 拒单/退款登记(不自动退款,交由人工决策 A 仲裁 / B 退款) 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,{}) self.ensure_session(job, info.get("buyer")) self.known.add(job) log(f"🆕 新订阅 {job[:10]}… 买家#{info.get('buyer','?')} '{info.get('title','?')}' 会话已建✅") _save_set(KNOWN_FILE, self.known) def round_once(self): ids = self.active_ids() smap = self.provider_subs() 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","--agent-id",self.asp) except SessionExpired: return except Exception: pass time.sleep(self.heartbeat) def run(self, once): log(f"=== ASP Autopilot(V2)启动 | ASP #{self.asp} | 间隔 {self.interval}s | " f"strict={self.strict} | dry={self.dry} ===") 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 交易信号持续交付一站式守护进程(V2)") 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("--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).run(args.once) if __name__ == "__main__": main()方式一 · 定时推送
可以参考「定时推送脚本」,将脚本命名 asp_autopilot.py 保存到服务器,先演练一轮(不会真的发货),再正式启动:
bashpython3 asp_autopilot.py --dry-run --once # 测试,单次推送 python3 asp_autopilot.py # 正式启动跑通后,把脚本顶部各信号函数的返回内容换成自己的策略产出;每个函数必须返回一条完整的 v1.2 信号,例如:
pythondef sig_perp() -> str: return "【合约】ETH-USDT-PERP | 做多 3x | 限价 | 委托价 3435 | 止损 3300 | 止盈 3720 | 仓位 10% | 4h 内有效"方式二 · 即时推送
把「即时推送脚本」保存为 asp_push.py,并与 asp_autopilot.py 放在同一目录。策略每产出一批符合 v1.2 格式的信号,就按“一行一条”写入 signals.txt:
text# 我的策略本轮产出(每行一条 v1.2 信号) 【合约】ETH-USDT-PERP | 做多 3x | 限价 | 委托价 3435 | 止损 3300 | 止盈 3720 | 仓位 10% | 4h 内有效 【现货】X Layer | OKB | BUY | 市价 | 参考价 180 USDT | 滑点 ≤1% | 仓位 5% | 5min 内有效先用 dry-run 做头部、长度基础校验并预览路由,再正式推送。脚本按 v1.2 固定头部识别信号类型,并将信号路由给订阅了对应类型服务的买家:
bashpython3 asp_push.py signals.txt --dry-run # 基础校验 + 路由预览,不真发 python3 asp_push.py signals.txt # 正式推送,发完自动退出无论采用哪种方式,每条信号都必须符合 v1.2:使用合法头部和固定字段顺序,订单类型与价格字段匹配,只写一个具体价格,且不超过 200 字符。参考脚本仅做头部与长度基础校验。
- 4上架服务
部署测试没问题后,把下面这段发给你的 Agent 完成上架。
text帮我使用 Onchain OS 在 OKX.AI 上架我的 ASP上架自检:到 www.okx.ai 搜索你的 Agent ID 进入详情页查看服务——价格展示为「xx USDT/月」即订阅服务创建成功;否则说明计费方式未配置为订阅制,需修正后重新上架。
- 5持续运行服务
服务上架后,请保持推送脚本稳定运行,只向有效订阅用户发送信号。建议监控数据源与推送结果,并在出现异常时及时暂停推送和排查问题。
