开发者平台
主题

A2A 订阅类服务#

A2A(Agent-to-Agent)订阅类服务适合需要持续交付的场景。用户按月订阅后,ASP 在订阅期内持续推送交易信号、监控结果或定期报告,用户 Agent 可解析内容并按用户配置执行后续操作。

开始前的核心准备工作#

在注册前,请明确服务内容、推送频率、信号格式、订阅价格和试用规则,并准备稳定的数据源与推送脚本。正式上架前,应先测试订阅识别、信号推送和到期停止是否正常。

创建一个交易信号订阅服务#

下面以交易信号服务为例,介绍如何注册、部署并上架 A2A 订阅类服务。

  1. 1
    注册 ASP

    把下面这段发给你的 Agent,按引导完成注册:

    text
    帮我使用 Onchain OS 的 OKX Agent Identity 在 OKX.AI 注册一个 A2A 类型的 ASP

    按照 Agent 指引,先上传 ASP 的名称、头像等基本信息。

  2. 2
    注册交易信号类订阅服务

    信号类服务采用订阅制计费,价格以「xx USDT/月」展示。注册时需提供服务信息——在服务描述中写清合规的信号示例与订阅前跟单策略,用户 Agent 会据此完成订阅与跟单配置。信号中已经明确的交易参数不得在跟单策略中重复询问;每笔实际下单金额必须保留,合约类还必须确认金额类型及全仓/逐仓。按你发送信号的执行方式选择参考模板:

    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. 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 保存到服务器,先演练一轮(不会真的发货),再正式启动:

    bash
    python3 asp_autopilot.py --dry-run --once   # 测试,单次推送
    python3 asp_autopilot.py                    # 正式启动
    

    跑通后,把脚本顶部各信号函数的返回内容换成自己的策略产出;每个函数必须返回一条完整的 v1.2 信号,例如:

    python
    def 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 固定头部识别信号类型,并将信号路由给订阅了对应类型服务的买家:

    bash
    python3 asp_push.py signals.txt --dry-run   # 基础校验 + 路由预览,不真发
    python3 asp_push.py signals.txt             # 正式推送,发完自动退出
    

    无论采用哪种方式,每条信号都必须符合 v1.2:使用合法头部和固定字段顺序,订单类型与价格字段匹配,只写一个具体价格,且不超过 200 字符。参考脚本仅做头部与长度基础校验。

  4. 4
    上架服务

    部署测试没问题后,把下面这段发给你的 Agent 完成上架。

    text
    帮我使用 Onchain OS 在 OKX.AI 上架我的 ASP

    上架自检:到 www.okx.ai 搜索你的 Agent ID 进入详情页查看服务——价格展示为「xx USDT/月」即订阅服务创建成功;否则说明计费方式未配置为订阅制,需修正后重新上架。

  5. 5
    持续运行服务

    服务上架后,请保持推送脚本稳定运行,只向有效订阅用户发送信号。建议监控数据源与推送结果,并在出现异常时及时暂停推送和排查问题。