连接 SDK

把应用接到 AIDC 的其他部分:企业数据集(固定报告)、实时数据流与智能体。

import { connect } from "/developer/sdk/v1/aidc.js";

数据集

const list = await connect.datasets.list();                         // 调用方能读的数据集
const report = await connect.datasets.report("aidc-demo-erp", "overview");

平台连接企业数据源,只出固定口径的聚合报告,不出明细行。应用只能读清单 datasets 里声明的数据集。

当前可用:

数据集 说明
aidc-demo-erp AIDC 合成 ERP 演示库(AIDC 数据库里的 demo_erp schema):虚构企业 12 个月的客户、供应商、物料、销售 / 采购订单、应收应付、收付款、库存流水。全部为虚构数据,给公开样板用

overview 报告包含:营收与毛利、订单与客单价、应收与逾期、DSO、库存金额、30 天内应付;月度趋势;客户 Top 10(含上下半年走势);品类结构;应收账龄;逾期客户;低于安全库存与呆滞物料;供应商采购额;以及规则化的经营信号(大客户下滑、逾期集中、营收异动、库存风险、回款周期)。

客户自己的 ERP(金蝶、用友、SAP、SQL Server…):连接在客户侧执行(数据不出客户环境),平台登记数据源、出同样形状的报告——接入由 AIDC 交付团队配置。

实时数据流

数据在源头(ERP、MES、SAP、OA…)一变,就由源头旁边的发布端推给 AIDC,AIDC 按序落账、保存最新状态,再在几秒内推给所有订阅端(Nexus 应用、CLI、智能体)。应用不轮询、不跑 cron,也不会有十几个看板各自去查同一个 ERP。

数据源 ──(发布端:aidc stream pipe + 适配器,装在数据源旁边)──HTTPS──▶ AIDC(按序落账 + 最新状态)──SSE──▶ 订阅端

订阅(在应用里)

清单里声明要订阅的流(与应用同一命名空间):"streams": ["orders"]。

const orders = connect.stream("orders");

const sub = orders.subscribe({
  onUpdate({ kind, rows, upserted, deleted, sourceAt }) {
    render(rows);                         // 合并好的当前全部行(按主键)
    flash(upserted);                      // 这次变化的行
  },
  onNotice({ data }) { toast(data.message); },            // 发布端发的通知(如「新合同」)
  onStatus({ connection, publisher }) {
    // connection: connecting / live / reconnecting / closed
    // publisher.live = false → 提示「数据可能不是最新」;publisher.status = "source_unreachable" → 「连不上数据源」
    badge(connection, publisher);
  },
});

const { rows, seq } = await orders.get();  // 一次性读当前状态

声明与发布

# 1. 声明:主键 + 字段白名单(发布的数据只许出现这些字段;证件号、银行卡、密码等类别直接拒绝)
aidc stream create orders --title "在手订单" --key order_no \
  --field order_no:string:合同号 --field qty:number:数量 --field due:datetime:交期

# 2. 签一把发布 Key:只能往这一条流发布,装在发布端(.env 0600)
aidc stream key orders --label "工厂服务器 · 订单发布端"

# 3. 在数据源旁边常驻运行发布端(systemd 服务,不是定时器)
AIDC_PUBLISH_KEY=aidc-pk-… aidc stream pipe cell-xxx/orders -- python3 orders_adapter.py

适配器是你写的常驻小程序(任何语言),用数据源自己的方式发现变化,往 stdout 每行打印一个 JSON:

行 含义
{"rows": [...]} 有主键的流的全表——发布端和上一次比较,只发变化的行(第一次发全量)
{"value": {...}} 没有主键的流的整体值——变了才发
{"event": {...}} 通知(字段同样受白名单约束,另可带 message)
{"status": "ok" | "degraded" | "source_unreachable", "detail": "…"} 适配器自报健康;每轮至少打一行,发布端 15 分钟收不到任何输出就报 degraded

aidc stream pipe 负责其余的事:算差、带幂等 id 重试(发布失败不丢、不重)、每 30 秒心跳、每 6 小时发一次全量校准、适配器崩溃按退避重启并上报「连不上数据源」。

# orders_adapter.py —— 常驻;连接保持;先用便宜的探针看有没有变化,变了才取全表
import json, time, pymssql

conn, last = None, None
while True:
    try:
        conn = conn or pymssql.connect(server="…", user="readonly", password=open("/path/secret").read().strip(), database="erp")
        cur = conn.cursor(as_dict=True)
        cur.execute("SELECT CHANGE_TRACKING_CURRENT_VERSION() AS v")   # 没开 Change Tracking 就用 rowversion / 更新时间水位
        version = cur.fetchone()["v"]
        if version != last:
            cur.execute("SELECT order_no, qty, due FROM dbo.orders WHERE status = 'open'")
            print(json.dumps({"rows": cur.fetchall()}, default=str), flush=True)
            last = version
        else:
            print('{"status": "ok"}', flush=True)
    except Exception as err:
        print(json.dumps({"status": "source_unreachable", "detail": str(err)[:200]}), flush=True)
        conn = None
    time.sleep(5)

也可以不用 CLI:在 Node / 服务端用 connect.stream("cell-xxx/orders").publisher()(rows() 自动算差、event()、status()),或直接 POST /api/v1/developer/streams/{namespace}/{name}/events。

怎么「发现变化」:按数据源

发布端之后的链路(落账、推送、续传、鉴权)对所有数据源都一样;不同的只是适配器怎么知道「变了」。能推送的源一律用推送;源头不能推送时,才在源头旁边做便宜的增量探针——一个探针服务所有订阅者,不是每个看板一个 cron。

数据源 最佳:源头推送 次选:边缘增量探针
SAP ECC / S/4HANA S/4 业务事件(Event Mesh)、IDoc / 变更指针、BAdI 调 HTTP RFC 按录入时间 / 计数器增量读(例:报工确认 AFRU 的 ERSDA/ERZET)
SQL Server(自研 / 行业 ERP) Change Tracking / CDC(DBA 开启一次) rowversion / 更新时间列水位;都没有时对「当天范围」做汇总指纹
PostgreSQL(MES 等) 主库逻辑复制 / LISTEN-NOTIFY 只读热备:WAL 回放位置门控 + 限频查询
金蝶 / 用友 / Odoo / SaaS ERP 开放平台 webhook、业务事件订阅 开放 API 按修改时间增量拉
OA(泛微等) 流程节点动作调 HTTP 流程表 requestid / 修改时间水位
企业微信 / 钉钉 事件回调 / Stream 模式 —
邮件 IMAP IDLE —

为什么不需要 Kafka:一条流的量级是每秒个位数事件、几个到几十个订阅者;要的是「有序、断线可补齐、新订阅者先拿最新状态、按公司鉴权」。AIDC 用数据库账本(有序、可补齐)+ 最新状态 + 已有的低延迟中继(唤醒订阅端)就够了,发布端只需要能出 HTTPS。到了每秒成千上万条事件、需要长时间重放给很多独立消费方时,再引入日志型中间件;即便那时,Debezium 这类 CDC 也可以直接把变化 POST 给同一个发布接口。

限制与权限

API

方法 路径 说明
GET / POST /api/v1/developer/streams 列出 / 声明(幂等)
GET / DELETE /api/v1/developer/streams/{namespace}/{name} 信息 + 最新状态 / 删除
GET /api/v1/developer/streams/{namespace}/{name}/events 订阅:Accept: text/event-stream 返回 SSE(state / change / notice / publisher / reconnect,每条 id = 序号,?after=N 续传);否则返回 N 之后的 JSON 增量
POST /api/v1/developer/streams/{namespace}/{name}/events 发布 {events:[{type, id?, at?, data}]}(type = snapshot / patch / event / status)
GET / POST / DELETE /api/v1/developer/streams/{namespace}/{name}/keys 发布 Key:列出 / 签发 / 吊销(?prefix=)

智能体

智能体走现成的 OpenAI 兼容 Agent API——与 Adis、外部框架调智能体是同一条路。凭证是访客自己的 Agent Key(aidc-sk-…,在 Console 的智能体页或 Adis 里获取);SDK 只在内存里使用,不存储、不上报。

const agent = connect.agent(agentKey);
const [info] = await agent.info();                    // 这把 Key 能调用的智能体
const reply = await agent.send("请把这份会议纪要归档,并跟进待办。\n\n" + minutes.markdown, { sessionId: "meeting-0925" });

for await (const delta of agent.stream("今天 3 号线有几件不合格?")) out.textContent += delta;

const fileId = await agent.upload(photo.blob, "defect.jpg");   // 图片 / 文件,24 小时内有效
await agent.send("请登记这件不合格品。", { files: [fileId] });

命令行

aidc connect datasets
aidc connect report aidc-demo-erp --json
AIDC_AGENT_KEY=aidc-sk-… aidc connect send "你好" --session demo

aidc stream create <流名> --title … --key <主键> --field 名:类型[:显示名] …
aidc stream key <流名> --label "装在哪"          # 发布 Key(aidc-pk-…),只显示一次
AIDC_PUBLISH_KEY=aidc-pk-… aidc stream pipe <命名空间>/<流名> -- <适配器命令…>
aidc stream tail <流名>                          # 实时订阅,每次变化一行
aidc stream status <流名> | get <流名> | list | keys | revoke | publish | delete

本页由 developer/docs/connect.md 生成 · Markdown 原文 · llms.txt