# 连接 SDK

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

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

## 数据集

```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"]`。

```js
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();  // 一次性读当前状态
```

- 第一次连上先给**全量**，之后每次变化给**增量**；断线自动重连并带着最后的序号续传——不丢、不重。
- 每条流都带发布端的健康状态：在线 / 离线（90 秒没心跳）、正常 / 部分异常 / 连不上数据源 / 已停止。数据源断了，订阅端**看得见**，不会像 cron 看板那样静默变旧。

### 声明与发布

```bash
# 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 小时发一次全量校准、适配器崩溃按退避重启并上报「连不上数据源」。

```python
# 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 给同一个发布接口。

### 限制与权限

- 单个事件 ≤ 1 MB、最新状态 ≤ 4 MB、每批 ≤ 100 条、每条流每分钟 ≤ 600 次发布；默认保留最近 1000 条事件供断线补齐（更早断开的订阅者直接拿最新状态）。
- 读：清单声明了这条流的同命名空间应用（客户应用只签给已登录的本公司成员）、本公司的开发者 Key；发布 Key 只能读、写它自己那一条流。
- 订阅连接约每 4 分 40 秒由服务端收尾一次，SDK 自动续接，应用无感。
- 删除数据流（`aidc stream delete`）会物理删除它的事件账与状态，并吊销它的发布 Key。

### 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 只在内存里使用，不存储、不上报。

```js
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] });
```

- `sessionId` 相同即同一段会话，智能体记得上文。
- 文件上传只支持云端运行的智能体；本地节点 / 快速车道的 Key 会报错。
- 在 AIDC 站点以外（CLI、服务端）用：`connect.agent(key, { agentBase: "https://api.ai-dc.ai/v1" })`。

## 命令行

```bash
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
```
