1052 lines
31 KiB
TypeScript
1052 lines
31 KiB
TypeScript
import { v7 as uuidv7 } from "uuid";
|
|
import {
|
|
APIError,
|
|
AuthReason,
|
|
Body,
|
|
ClientOptions,
|
|
ConnectionEvent,
|
|
ConnectionState,
|
|
Credential,
|
|
GroupEvent,
|
|
GroupMemberIn,
|
|
HandshakeLimits,
|
|
Message,
|
|
PresenceEvent,
|
|
RecallResult,
|
|
Receipt,
|
|
RegisterOptions,
|
|
RegisterResult,
|
|
ReconnectBackoff,
|
|
RevokedEvent,
|
|
SendOptions,
|
|
SendResult,
|
|
Target,
|
|
Transport,
|
|
applyJitter,
|
|
marshalJSON,
|
|
nominalDelay,
|
|
registerURLFromConnect,
|
|
} from "./types.js";
|
|
import { MqttTransport } from "./mqtt.js";
|
|
|
|
type RespFrame = {
|
|
ok: boolean;
|
|
data?: unknown;
|
|
error?: { code: string; message: string };
|
|
};
|
|
|
|
type Pending = {
|
|
resolve: (r: RespFrame) => void;
|
|
reject: (e: unknown) => void;
|
|
isSend?: boolean;
|
|
timer?: ReturnType<typeof setTimeout>;
|
|
};
|
|
|
|
type SendItem = {
|
|
frame: Record<string, unknown>;
|
|
payload: string;
|
|
id: string;
|
|
result: { resolve: (v: SendResult) => void; reject: (e: unknown) => void };
|
|
inflight: boolean;
|
|
epoch: number;
|
|
rateN: number;
|
|
abandoned?: boolean;
|
|
abort?: () => void;
|
|
};
|
|
|
|
type DedupState = "delivered" | "acked" | "revoked";
|
|
|
|
class LRUMap<T> {
|
|
private map = new Map<string, T>();
|
|
constructor(private cap: number) {}
|
|
get(k: string): T | undefined {
|
|
const v = this.map.get(k);
|
|
if (v !== undefined) {
|
|
this.map.delete(k);
|
|
this.map.set(k, v);
|
|
}
|
|
return v;
|
|
}
|
|
has(k: string): boolean {
|
|
return this.map.has(k);
|
|
}
|
|
put(k: string, v: T): void {
|
|
if (this.map.has(k)) this.map.delete(k);
|
|
this.map.set(k, v);
|
|
while (this.map.size > this.cap) {
|
|
const first = this.map.keys().next().value as string | undefined;
|
|
if (first === undefined) break;
|
|
this.map.delete(first);
|
|
}
|
|
}
|
|
delete(k: string): void {
|
|
this.map.delete(k);
|
|
}
|
|
}
|
|
|
|
export class Client {
|
|
private opts: Required<
|
|
Pick<
|
|
ClientOptions,
|
|
| "manualAck"
|
|
| "allowTcp"
|
|
| "connectTimeoutMs"
|
|
| "clientLabel"
|
|
| "sendQueueSize"
|
|
| "maxInflight"
|
|
| "dedupCapacity"
|
|
>
|
|
> &
|
|
ClientOptions = {
|
|
manualAck: false,
|
|
allowTcp: false,
|
|
connectTimeoutMs: 30_000,
|
|
clientLabel: "js-sdk/0.1",
|
|
sendQueueSize: 1000,
|
|
maxInflight: 100,
|
|
dedupCapacity: 10000,
|
|
};
|
|
|
|
private transport?: Transport;
|
|
private backoff = new ReconnectBackoff();
|
|
private endpointId = "";
|
|
private ridSeq = 0;
|
|
private pending = new Map<string, Pending>();
|
|
private sendQ: SendItem[] = [];
|
|
private inflight = 0;
|
|
private handshook = false;
|
|
private stopReconnect = false;
|
|
private closed = false;
|
|
private lastStopCode = "";
|
|
private lastStopErr?: APIError;
|
|
private state: ConnectionState = "offline";
|
|
private limits: HandshakeLimits = {
|
|
server_time_ms: 0,
|
|
server_version: "",
|
|
max_body_bytes: 262144,
|
|
max_meta_bytes: 4096,
|
|
max_frame_bytes: 786432,
|
|
max_ttl_seconds: 2592000,
|
|
max_schedule_seconds: 31536000,
|
|
ack_timeout_seconds: 300,
|
|
};
|
|
private clockSkew = 0;
|
|
private store = new LRUMap<DedupState>(10000);
|
|
private cbChain: Promise<void> = Promise.resolve();
|
|
private timers = new Set<ReturnType<typeof setTimeout>>();
|
|
private watch?: { ids?: string[]; all?: boolean };
|
|
|
|
private onSession?: (token: string) => void;
|
|
private onMessage?: (msg: Message) => void | Promise<void>;
|
|
private onReceipt?: (r: Receipt) => void;
|
|
private onRevoked?: (e: RevokedEvent) => void;
|
|
private onPresence?: (e: PresenceEvent) => void;
|
|
private onGroupEvent?: (e: GroupEvent) => void;
|
|
private onConnection?: (e: ConnectionEvent) => void;
|
|
|
|
onSessionHandler(h: (token: string) => void): void {
|
|
this.onSession = h;
|
|
}
|
|
onMessageHandler(h: (msg: Message) => void | Promise<void>): void {
|
|
this.onMessage = h;
|
|
}
|
|
onReceiptHandler(h: (r: Receipt) => void): void {
|
|
this.onReceipt = h;
|
|
}
|
|
onRevokedHandler(h: (e: RevokedEvent) => void): void {
|
|
this.onRevoked = h;
|
|
}
|
|
onPresenceHandler(h: (e: PresenceEvent) => void): void {
|
|
this.onPresence = h;
|
|
}
|
|
onGroupEventHandler(h: (e: GroupEvent) => void): void {
|
|
this.onGroupEvent = h;
|
|
}
|
|
onConnectionHandler(h: (e: ConnectionEvent) => void): void {
|
|
this.onConnection = h;
|
|
}
|
|
|
|
onSessionCb(h: (token: string) => void): void {
|
|
this.onSessionHandler(h);
|
|
}
|
|
|
|
clockSkewMs(): number {
|
|
return this.clockSkew;
|
|
}
|
|
|
|
getLimits(): HandshakeLimits {
|
|
return { ...this.limits };
|
|
}
|
|
|
|
lastStopCodeForTest(): string {
|
|
return this.lastStopCode;
|
|
}
|
|
|
|
resendPayloadForTest(): string | undefined {
|
|
return this.sendQ[0]?.payload;
|
|
}
|
|
|
|
async connect(
|
|
url: string,
|
|
endpointId: string,
|
|
credential: Credential,
|
|
options: ClientOptions = {},
|
|
): Promise<void> {
|
|
if (this.closed) throw new APIError("closed", "已关闭");
|
|
if (this.transport) throw new APIError("bad_request", "已在连接中");
|
|
if (options.maxReceiveBytes != null && options.maxReceiveBytes > 0 && options.maxReceiveBytes < 1024) {
|
|
throw new APIError("bad_request", "max_receive_bytes 小于 1024");
|
|
}
|
|
this.opts = {
|
|
...this.opts,
|
|
...options,
|
|
connectTimeoutMs: options.connectTimeoutMs ?? 30_000,
|
|
clientLabel: options.clientLabel ?? "js-sdk/0.1",
|
|
sendQueueSize: options.sendQueueSize ?? 1000,
|
|
maxInflight: options.maxInflight ?? 100,
|
|
dedupCapacity: options.dedupCapacity ?? 10000,
|
|
manualAck: options.manualAck ?? false,
|
|
allowTcp: options.allowTcp ?? false,
|
|
};
|
|
this.store = new LRUMap<DedupState>(this.opts.dedupCapacity!);
|
|
this.endpointId = endpointId;
|
|
this.stopReconnect = false;
|
|
this.handshook = false;
|
|
this.lastStopCode = "";
|
|
this.lastStopErr = undefined;
|
|
this.backoff = new ReconnectBackoff();
|
|
const pass = credential.sessionToken ?? credential.password ?? "";
|
|
const tr = options.transport ?? new MqttTransport();
|
|
this.transport = tr;
|
|
tr.setCredential(pass);
|
|
this.setState("connecting");
|
|
|
|
void tr
|
|
.start({
|
|
url,
|
|
endpointId,
|
|
connectTimeoutMs: this.opts.connectTimeoutMs!,
|
|
allowTcp: !!this.opts.allowTcp,
|
|
backoff: this.backoff,
|
|
onDown: (p) => this.handleDown(p),
|
|
onOffline: () => this.handleOffline(),
|
|
onAuthFailed: (r) => this.failAuth(r),
|
|
onKicked: () => this.failKicked(),
|
|
mqttReady: () => this.doHello(),
|
|
})
|
|
.catch(() => {});
|
|
|
|
const deadline = Date.now() + this.opts.connectTimeoutMs!;
|
|
while (Date.now() < deadline) {
|
|
if (this.handshook) return;
|
|
if (this.stopReconnect || this.state === "auth_failed" || this.state === "kicked") {
|
|
throw this.lastStopErr ?? new APIError(this.lastStopCode || "auth_failed", this.state);
|
|
}
|
|
await sleep(20);
|
|
}
|
|
await this.teardown(false, new APIError("not_connected", "连接超时"));
|
|
throw new APIError("not_connected", "连接超时");
|
|
}
|
|
|
|
private failAuth(reason: AuthReason): void {
|
|
if (this.stopReconnect && this.lastStopCode) return;
|
|
const err = new APIError(reason, "认证失败,停止重连");
|
|
this.stopReconnect = true;
|
|
this.handshook = false;
|
|
this.lastStopCode = reason;
|
|
this.lastStopErr = err;
|
|
this.setState("auth_failed", reason);
|
|
this.failQueued(err);
|
|
this.failPending(err, true);
|
|
void this.transport?.stop();
|
|
this.transport = undefined;
|
|
}
|
|
|
|
private failKicked(): void {
|
|
if (this.stopReconnect && this.lastStopCode === "taken_over") return;
|
|
const err = new APIError("taken_over", "被顶号,停止重连");
|
|
this.stopReconnect = true;
|
|
this.handshook = false;
|
|
this.lastStopCode = "taken_over";
|
|
this.lastStopErr = err;
|
|
this.setState("kicked", "taken_over");
|
|
this.failQueued(err);
|
|
this.failPending(err, true);
|
|
void this.transport?.stop();
|
|
this.transport = undefined;
|
|
}
|
|
|
|
private handleFatal(reason: string): void {
|
|
if (this.stopReconnect && this.lastStopCode) return;
|
|
const code = reason || "fatal";
|
|
const err = new APIError(code, "致命错误,停止重连");
|
|
this.stopReconnect = true;
|
|
this.handshook = false;
|
|
this.lastStopCode = code;
|
|
this.lastStopErr = err;
|
|
this.setState("auth_failed", code);
|
|
this.failQueued(err);
|
|
this.failPending(err, true);
|
|
void this.transport?.stop();
|
|
this.transport = undefined;
|
|
}
|
|
|
|
private handleOffline(): void {
|
|
this.handshook = false;
|
|
if (this.stopReconnect || this.closed) {
|
|
this.failPending(new APIError("not_connected", "未连接"), true);
|
|
return;
|
|
}
|
|
this.requeueInflight();
|
|
this.failPending(new APIError("not_connected", "未连接"), false);
|
|
this.setState("reconnecting");
|
|
}
|
|
|
|
private requeueInflight(): void {
|
|
for (const it of this.sendQ) {
|
|
if (!it.inflight) continue;
|
|
it.epoch++;
|
|
it.inflight = false;
|
|
const rid = String(it.frame.rid ?? "");
|
|
const p = this.pending.get(rid);
|
|
if (p) {
|
|
this.clearPendingTimer(p);
|
|
this.pending.delete(rid);
|
|
}
|
|
this.regenerateSend(it);
|
|
}
|
|
this.inflight = 0;
|
|
}
|
|
|
|
private regenerateSend(it: SendItem): void {
|
|
it.frame.rid = this.nextRid();
|
|
it.payload = marshalJSON(it.frame);
|
|
}
|
|
|
|
private setState(state: ConnectionState, reason?: string): void {
|
|
this.state = state;
|
|
this.enqueueCb(() => this.onConnection?.({ state, reason }));
|
|
}
|
|
|
|
private enqueueCb(fn: () => void | Promise<void>): void {
|
|
this.cbChain = this.cbChain.then(async () => {
|
|
try {
|
|
await fn();
|
|
} catch {
|
|
/* 回调错误不打断串行链 */
|
|
}
|
|
});
|
|
}
|
|
|
|
private nextRid(): string {
|
|
this.ridSeq += 1;
|
|
return String(this.ridSeq);
|
|
}
|
|
|
|
private stopErr(): APIError {
|
|
if (this.lastStopErr) return this.lastStopErr;
|
|
if (this.lastStopCode) return new APIError(this.lastStopCode);
|
|
if (this.closed) return new APIError("closed", "已关闭");
|
|
return new APIError("not_connected");
|
|
}
|
|
|
|
private async doHello(): Promise<void> {
|
|
const sentAt = Date.now();
|
|
const req: Record<string, unknown> = {
|
|
v: 1,
|
|
type: "hello",
|
|
rid: this.nextRid(),
|
|
client: this.opts.clientLabel,
|
|
};
|
|
if (this.opts.maxReceiveBytes && this.opts.maxReceiveBytes > 0) {
|
|
req.max_receive_bytes = this.opts.maxReceiveBytes;
|
|
}
|
|
const data = (await this.request(req, true)) as Record<string, unknown>;
|
|
const recvAt = Date.now();
|
|
const serverTime = Number(data.server_time_ms ?? 0);
|
|
this.clockSkew = serverTime - Math.floor((sentAt + recvAt) / 2);
|
|
this.limits = {
|
|
server_time_ms: serverTime,
|
|
server_version: String(data.server_version ?? ""),
|
|
max_body_bytes: Number(data.max_body_bytes ?? 262144),
|
|
max_meta_bytes: Number(data.max_meta_bytes ?? 4096),
|
|
max_frame_bytes: Number(data.max_frame_bytes ?? 786432),
|
|
max_ttl_seconds: Number(data.max_ttl_seconds ?? 2592000),
|
|
max_schedule_seconds: Number(data.max_schedule_seconds ?? 31536000),
|
|
ack_timeout_seconds: Number(data.ack_timeout_seconds ?? 300),
|
|
};
|
|
this.handshook = true;
|
|
this.setState("online");
|
|
const token = data.session_token ? String(data.session_token) : "";
|
|
if (token) {
|
|
this.transport?.setCredential(token);
|
|
this.enqueueCb(() => this.onSession?.(token));
|
|
}
|
|
if (this.watch) {
|
|
void this.restoreWatch();
|
|
}
|
|
void this.drainSendQueue();
|
|
}
|
|
|
|
private async restoreWatch(): Promise<void> {
|
|
if (!this.watch) return;
|
|
const req: Record<string, unknown> = { v: 1, type: "presence.watch", rid: this.nextRid() };
|
|
if (this.watch.all) req.all = true;
|
|
else req.ids = this.watch.ids ?? [];
|
|
try {
|
|
await this.request(req, false);
|
|
} catch {
|
|
/* 重连后尽力恢复 */
|
|
}
|
|
}
|
|
|
|
private handleDown(payload: Uint8Array): void {
|
|
let head: { type?: string; rid?: string };
|
|
try {
|
|
head = JSON.parse(new TextDecoder().decode(payload));
|
|
} catch {
|
|
return;
|
|
}
|
|
const text = new TextDecoder().decode(payload);
|
|
switch (head.type) {
|
|
case "resp": {
|
|
const rf = JSON.parse(text) as RespFrame & { rid: string };
|
|
const p = this.pending.get(rf.rid ?? head.rid!);
|
|
if (p) {
|
|
this.clearPendingTimer(p);
|
|
this.pending.delete(rf.rid ?? head.rid!);
|
|
p.resolve(rf);
|
|
}
|
|
break;
|
|
}
|
|
case "msg":
|
|
void this.handleMsg(JSON.parse(text)).catch(() => {});
|
|
break;
|
|
case "receipt":
|
|
void this.handleReceipt(JSON.parse(text)).catch(() => {});
|
|
break;
|
|
case "revoked":
|
|
this.handleRevoked(JSON.parse(text));
|
|
break;
|
|
case "presence": {
|
|
const p = JSON.parse(text) as PresenceEvent;
|
|
this.enqueueCb(() => this.onPresence?.(p));
|
|
break;
|
|
}
|
|
case "group_event": {
|
|
const g = JSON.parse(text) as GroupEvent;
|
|
this.enqueueCb(() => this.onGroupEvent?.(g));
|
|
break;
|
|
}
|
|
case "fatal": {
|
|
const f = JSON.parse(text) as { reason?: string };
|
|
this.handleFatal(f.reason ?? "fatal");
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
private async handleMsg(m: {
|
|
id: string;
|
|
from: string;
|
|
to: Target;
|
|
body: Body;
|
|
meta?: Record<string, unknown>;
|
|
send_at_ms: number;
|
|
}): Promise<void> {
|
|
try {
|
|
const key = `m\0${m.from}\0${m.id}`;
|
|
const ent = this.store.get(key);
|
|
if (ent === "acked") {
|
|
void this.sendAckFrame(m.from, m.id).catch(() => {});
|
|
return;
|
|
}
|
|
if (ent === "delivered" || ent === "revoked") return;
|
|
this.store.put(key, "delivered");
|
|
|
|
const msg: Message = {
|
|
id: m.id,
|
|
from: m.from,
|
|
to: m.to,
|
|
body: m.body,
|
|
meta: m.meta,
|
|
send_at_ms: m.send_at_ms,
|
|
};
|
|
|
|
let cbErr: unknown;
|
|
await new Promise<void>((resolve) => {
|
|
this.enqueueCb(async () => {
|
|
if (this.store.get(key) === "revoked") {
|
|
resolve();
|
|
return;
|
|
}
|
|
try {
|
|
await this.onMessage?.(msg);
|
|
} catch (e) {
|
|
cbErr = e;
|
|
}
|
|
resolve();
|
|
});
|
|
});
|
|
|
|
if (this.store.get(key) === "revoked") return;
|
|
if (this.opts.manualAck) return;
|
|
if (cbErr) {
|
|
this.store.delete(key);
|
|
return;
|
|
}
|
|
// S-06:回调成功即置 acked,再发 ack
|
|
this.store.put(key, "acked");
|
|
void this.sendAckFrame(m.from, m.id).catch(() => {});
|
|
} catch {
|
|
/* S-04:确认路径不得产生未处理拒绝 */
|
|
}
|
|
}
|
|
|
|
async ack(msg: Message): Promise<void> {
|
|
await this.sendAckFrame(msg.from, msg.id);
|
|
this.store.put(`m\0${msg.from}\0${msg.id}`, "acked");
|
|
}
|
|
|
|
private async sendAckFrame(from: string, id: string): Promise<void> {
|
|
const data = await this.request(
|
|
{ v: 1, type: "ack", rid: this.nextRid(), from, id },
|
|
true,
|
|
);
|
|
if (data && typeof data === "object" && "result" in (data as object)) {
|
|
const result = String((data as { result: string }).result);
|
|
if (result && result !== "accepted") {
|
|
this.enqueueCb(() => this.onRevoked?.({ id, from, reason: result }));
|
|
}
|
|
}
|
|
}
|
|
|
|
private async handleReceipt(r: Receipt & { receipt_id: string }): Promise<void> {
|
|
try {
|
|
const key = `r\0${r.receipt_id}`;
|
|
if (this.store.has(key)) {
|
|
void this.sendReceiptAck(r.receipt_id).catch(() => {});
|
|
return;
|
|
}
|
|
this.store.put(key, "acked");
|
|
this.enqueueCb(() => this.onReceipt?.(r));
|
|
void this.sendReceiptAck(r.receipt_id).catch(() => {});
|
|
} catch {
|
|
/* S-04 */
|
|
}
|
|
}
|
|
|
|
private async sendReceiptAck(receiptId: string): Promise<void> {
|
|
await this.request({ v: 1, type: "receipt_ack", rid: this.nextRid(), receipt_id: receiptId }, true);
|
|
}
|
|
|
|
private handleRevoked(r: RevokedEvent): void {
|
|
const key = `m\0${r.from}\0${r.id}`;
|
|
const ent = this.store.get(key);
|
|
if (ent === "acked" || ent === "revoked") return;
|
|
this.store.put(key, "revoked");
|
|
this.enqueueCb(() => this.onRevoked?.(r));
|
|
}
|
|
|
|
private request(frame: Record<string, unknown>, allowUnready: boolean): Promise<unknown> {
|
|
if (this.closed) return Promise.reject(new APIError("closed"));
|
|
if (!allowUnready && !this.handshook) return Promise.reject(new APIError("not_connected", "未握手"));
|
|
const tr = this.transport;
|
|
if (!tr) return Promise.reject(new APIError("not_connected"));
|
|
const rid = String(frame.rid ?? this.nextRid());
|
|
frame.rid = rid;
|
|
const payload = marshalJSON(frame);
|
|
return new Promise((resolve, reject) => {
|
|
const p: Pending = {
|
|
resolve: (rf) => {
|
|
this.clearPendingTimer(p);
|
|
if (!rf.ok) {
|
|
reject(new APIError(rf.error?.code ?? "bad_request", rf.error?.message ?? ""));
|
|
return;
|
|
}
|
|
resolve(rf.data);
|
|
},
|
|
reject: (e) => {
|
|
this.clearPendingTimer(p);
|
|
reject(e);
|
|
},
|
|
isSend: false,
|
|
};
|
|
p.timer = setTimeout(() => {
|
|
this.pending.delete(rid);
|
|
p.reject(new APIError("not_connected", "请求超时"));
|
|
}, 60_000);
|
|
p.timer.unref?.();
|
|
this.timers.add(p.timer);
|
|
this.pending.set(rid, p);
|
|
void tr.publishUp(payload).catch((e) => {
|
|
this.pending.delete(rid);
|
|
this.clearPendingTimer(p);
|
|
reject(e instanceof APIError ? e : new APIError("not_connected", String(e)));
|
|
});
|
|
});
|
|
}
|
|
|
|
async send(to: Target, body: Body, opt: SendOptions = {}): Promise<SendResult> {
|
|
if (this.closed || this.stopReconnect) throw this.stopErr();
|
|
|
|
const enc = body.enc || "utf8";
|
|
const b: Body = {
|
|
enc,
|
|
data: body.data,
|
|
content_type:
|
|
opt.contentType ||
|
|
body.content_type ||
|
|
(enc === "base64" ? "application/octet-stream" : "text/plain; charset=utf-8"),
|
|
};
|
|
const n = bodyDecodedLen(b);
|
|
const maxBody = this.limits.max_body_bytes || 262144;
|
|
if (n > maxBody) throw new APIError("body_too_large", "正文超限");
|
|
|
|
if (opt.sendAt && opt.delayMs != null) throw new APIError("bad_request", "sendAt 与 delay 互斥");
|
|
|
|
const id = opt.id || uuidv7();
|
|
const frame: Record<string, unknown> = {
|
|
v: 1,
|
|
type: "send",
|
|
rid: this.nextRid(),
|
|
id,
|
|
to,
|
|
body: b,
|
|
};
|
|
if (opt.meta) frame.meta = opt.meta;
|
|
if (opt.talkPassword) frame.talk_password = opt.talkPassword;
|
|
if (opt.receipt != null) frame.receipt = opt.receipt;
|
|
if (opt.keep) {
|
|
const off: Record<string, unknown> = { keep: true };
|
|
if (opt.ttl != null) off.ttl_seconds = opt.ttl;
|
|
frame.offline = off;
|
|
}
|
|
if (opt.sendAt) {
|
|
frame.send_at_ms = opt.sendAt.getTime() + this.clockSkew;
|
|
} else if (opt.delayMs != null) {
|
|
frame.delay_ms = opt.delayMs;
|
|
}
|
|
|
|
const payload = marshalJSON(frame);
|
|
const maxFrame = this.limits.max_frame_bytes || 786432;
|
|
if (utf8Len(payload) > maxFrame) {
|
|
throw new APIError("frame_too_large", "整帧超限");
|
|
}
|
|
if (this.sendQ.length >= this.opts.sendQueueSize!) {
|
|
throw new APIError("queue_full", "发送队列已满");
|
|
}
|
|
|
|
return new Promise<SendResult>((resolve, reject) => {
|
|
const item: SendItem = {
|
|
frame,
|
|
payload,
|
|
id,
|
|
result: { resolve, reject },
|
|
inflight: false,
|
|
epoch: 0,
|
|
rateN: 0,
|
|
};
|
|
if (opt.signal) {
|
|
const onAbort = () => {
|
|
if (item.inflight) {
|
|
item.abandoned = true;
|
|
reject(new APIError("result_unknown", "结果未知,请用同一消息号重试"));
|
|
return;
|
|
}
|
|
this.sendQ = this.sendQ.filter((x) => x !== item);
|
|
reject(new APIError("closed", "已取消"));
|
|
};
|
|
if (opt.signal.aborted) {
|
|
onAbort();
|
|
return;
|
|
}
|
|
opt.signal.addEventListener("abort", onAbort, { once: true });
|
|
item.abort = () => opt.signal?.removeEventListener("abort", onAbort);
|
|
}
|
|
this.sendQ.push(item);
|
|
void this.drainSendQueue();
|
|
});
|
|
}
|
|
|
|
private async drainSendQueue(): Promise<void> {
|
|
while (true) {
|
|
if (!this.handshook || !this.transport || this.stopReconnect) return;
|
|
const maxFrame = this.limits.max_frame_bytes || 786432;
|
|
const next = this.sendQ.find((x) => !x.inflight && !x.abandoned);
|
|
if (!next || this.inflight >= this.opts.maxInflight!) return;
|
|
if (utf8Len(next.payload) > maxFrame) {
|
|
this.finishSendErr(next, new APIError("frame_too_large", "整帧超限"));
|
|
continue;
|
|
}
|
|
next.epoch++;
|
|
next.inflight = true;
|
|
this.inflight++;
|
|
void this.dispatchSend(next, next.epoch);
|
|
}
|
|
}
|
|
|
|
private async dispatchSend(item: SendItem, epoch: number): Promise<void> {
|
|
const rid = String(item.frame.rid);
|
|
const tr = this.transport;
|
|
if (!tr) return;
|
|
try {
|
|
const data = await new Promise<unknown>((resolve, reject) => {
|
|
if (item.epoch !== epoch || !item.inflight) {
|
|
reject(new Error("stale"));
|
|
return;
|
|
}
|
|
const p: Pending = {
|
|
resolve: (rf) => {
|
|
this.clearPendingTimer(p);
|
|
if (!rf.ok) {
|
|
reject(new APIError(rf.error?.code ?? "bad_request", rf.error?.message ?? ""));
|
|
return;
|
|
}
|
|
resolve(rf.data);
|
|
},
|
|
reject: (e) => {
|
|
this.clearPendingTimer(p);
|
|
reject(e);
|
|
},
|
|
isSend: true,
|
|
};
|
|
this.pending.set(rid, p);
|
|
void tr.publishUp(item.payload).catch((e) => {
|
|
this.pending.delete(rid);
|
|
this.clearPendingTimer(p);
|
|
reject(e);
|
|
});
|
|
});
|
|
if (item.epoch !== epoch) return;
|
|
const sd = (data ?? {}) as SendResult;
|
|
this.finishSend(item, { id: sd.id || item.id, send_at_ms: sd.send_at_ms, state: sd.state });
|
|
} catch (e) {
|
|
if (item.epoch !== epoch) return;
|
|
if (e instanceof APIError && e.code === "rate_limited") {
|
|
item.inflight = false;
|
|
this.inflight = Math.max(0, this.inflight - 1);
|
|
this.pending.delete(rid);
|
|
item.rateN++;
|
|
this.regenerateSend(item);
|
|
const wait = applyJitter(nominalDelay(item.rateN));
|
|
this.schedule(() => void this.drainSendQueue(), wait);
|
|
return;
|
|
}
|
|
if (e instanceof APIError && e.code === "_retry") return;
|
|
if (!(e instanceof APIError)) {
|
|
if (this.stopReconnect) {
|
|
this.finishSendErr(item, this.stopErr());
|
|
return;
|
|
}
|
|
item.inflight = false;
|
|
this.inflight = Math.max(0, this.inflight - 1);
|
|
this.pending.delete(rid);
|
|
this.regenerateSend(item);
|
|
return;
|
|
}
|
|
this.finishSendErr(item, e);
|
|
}
|
|
}
|
|
|
|
private finishSend(item: SendItem, res: SendResult): void {
|
|
this.sendQ = this.sendQ.filter((x) => x !== item);
|
|
if (item.inflight) {
|
|
this.inflight = Math.max(0, this.inflight - 1);
|
|
item.inflight = false;
|
|
}
|
|
const rid = String(item.frame.rid ?? "");
|
|
this.pending.delete(rid);
|
|
item.abort?.();
|
|
if (!item.abandoned) item.result.resolve(res);
|
|
void this.drainSendQueue();
|
|
}
|
|
|
|
private finishSendErr(item: SendItem, err: unknown): void {
|
|
this.sendQ = this.sendQ.filter((x) => x !== item);
|
|
if (item.inflight) {
|
|
this.inflight = Math.max(0, this.inflight - 1);
|
|
item.inflight = false;
|
|
}
|
|
const rid = String(item.frame.rid ?? "");
|
|
this.pending.delete(rid);
|
|
item.abort?.();
|
|
if (!item.abandoned) item.result.reject(err);
|
|
void this.drainSendQueue();
|
|
}
|
|
|
|
private failQueued(err: unknown): void {
|
|
for (const it of this.sendQ) {
|
|
it.abort?.();
|
|
if (!it.abandoned) it.result.reject(err);
|
|
}
|
|
this.sendQ = [];
|
|
this.inflight = 0;
|
|
}
|
|
|
|
private failPending(err: unknown, all: boolean): void {
|
|
for (const [rid, p] of [...this.pending]) {
|
|
this.clearPendingTimer(p);
|
|
this.pending.delete(rid);
|
|
if (!all && p.isSend) {
|
|
p.reject(new APIError("_retry", "disconnected"));
|
|
continue;
|
|
}
|
|
p.reject(err);
|
|
}
|
|
}
|
|
|
|
private clearPendingTimer(p: Pending): void {
|
|
if (p.timer) {
|
|
clearTimeout(p.timer);
|
|
this.timers.delete(p.timer);
|
|
p.timer = undefined;
|
|
}
|
|
}
|
|
|
|
private schedule(fn: () => void, ms: number): void {
|
|
const t = setTimeout(() => {
|
|
this.timers.delete(t);
|
|
fn();
|
|
}, ms);
|
|
t.unref?.();
|
|
this.timers.add(t);
|
|
}
|
|
|
|
private clearTimers(): void {
|
|
for (const t of this.timers) clearTimeout(t);
|
|
this.timers.clear();
|
|
}
|
|
|
|
private async teardown(setClosed: boolean, stopErr: APIError): Promise<void> {
|
|
if (setClosed) this.closed = true;
|
|
this.stopReconnect = true;
|
|
this.lastStopErr = stopErr;
|
|
this.lastStopCode = stopErr.code;
|
|
this.failQueued(stopErr);
|
|
this.failPending(stopErr, true);
|
|
this.handshook = false;
|
|
this.clearTimers();
|
|
this.setState("offline", this.lastStopCode);
|
|
const tr = this.transport;
|
|
this.transport = undefined;
|
|
await tr?.stop();
|
|
}
|
|
|
|
async recall(id: string): Promise<RecallResult> {
|
|
return (await this.request({ v: 1, type: "recall", rid: this.nextRid(), id }, false)) as RecallResult;
|
|
}
|
|
|
|
async status(id: string, cursor = "", limit = 0): Promise<unknown> {
|
|
const req: Record<string, unknown> = { v: 1, type: "status", rid: this.nextRid(), id };
|
|
if (cursor) req.cursor = cursor;
|
|
if (limit) req.limit = limit;
|
|
return this.request(req, false);
|
|
}
|
|
|
|
async unlock(endpointId: string, talkPassword: string): Promise<void> {
|
|
await this.request(
|
|
{ v: 1, type: "unlock", rid: this.nextRid(), endpoint_id: endpointId, talk_password: talkPassword },
|
|
false,
|
|
);
|
|
}
|
|
|
|
async presence(ids: string[]): Promise<unknown> {
|
|
return this.request({ v: 1, type: "presence.get", rid: this.nextRid(), ids }, false);
|
|
}
|
|
|
|
async directory(cursor = "", query = "", limit = 0): Promise<unknown> {
|
|
const req: Record<string, unknown> = { v: 1, type: "directory.list", rid: this.nextRid() };
|
|
if (cursor) req.cursor = cursor;
|
|
if (query) req.query = query;
|
|
if (limit) req.limit = limit;
|
|
return this.request(req, false);
|
|
}
|
|
|
|
async watchPresence(ids: string[] | "all"): Promise<void> {
|
|
if (ids === "all") this.watch = { all: true };
|
|
else this.watch = { ids: [...ids] };
|
|
const req: Record<string, unknown> = { v: 1, type: "presence.watch", rid: this.nextRid() };
|
|
if (ids === "all") req.all = true;
|
|
else req.ids = ids;
|
|
await this.request(req, false);
|
|
}
|
|
|
|
async getSelf(): Promise<unknown> {
|
|
return this.request({ v: 1, type: "self.get", rid: this.nextRid() }, false);
|
|
}
|
|
|
|
async updateSelf(name?: string, defaultDelayMs?: number): Promise<void> {
|
|
const req: Record<string, unknown> = { v: 1, type: "self.update", rid: this.nextRid() };
|
|
if (name != null) req.name = name;
|
|
if (defaultDelayMs != null) req.default_delay_ms = defaultDelayMs;
|
|
await this.request(req, false);
|
|
}
|
|
|
|
async setTalkPassword(talkPassword: string): Promise<void> {
|
|
await this.request(
|
|
{ v: 1, type: "self.talk_password", rid: this.nextRid(), talk_password: talkPassword },
|
|
false,
|
|
);
|
|
}
|
|
|
|
async changeLoginPassword(oldPassword: string, newPassword: string): Promise<void> {
|
|
const data = (await this.request(
|
|
{
|
|
v: 1,
|
|
type: "self.login_password",
|
|
rid: this.nextRid(),
|
|
old_password: oldPassword,
|
|
new_password: newPassword,
|
|
},
|
|
false,
|
|
)) as { session_token?: string };
|
|
if (data?.session_token) {
|
|
this.transport?.setCredential(data.session_token);
|
|
this.enqueueCb(() => this.onSession?.(data.session_token!));
|
|
}
|
|
}
|
|
|
|
async createGroup(id: string, name: string, members: GroupMemberIn[]): Promise<unknown> {
|
|
return this.request(
|
|
{
|
|
v: 1,
|
|
type: "group.create",
|
|
rid: this.nextRid(),
|
|
id,
|
|
name,
|
|
members: members.map((m) => ({ id: m.id, talk_password: m.talkPassword ?? "" })),
|
|
},
|
|
false,
|
|
);
|
|
}
|
|
|
|
async addGroupMembers(groupId: string, members: GroupMemberIn[]): Promise<unknown> {
|
|
return this.request(
|
|
{
|
|
v: 1,
|
|
type: "group.add",
|
|
rid: this.nextRid(),
|
|
group_id: groupId,
|
|
members: members.map((m) => ({ id: m.id, talk_password: m.talkPassword ?? "" })),
|
|
},
|
|
false,
|
|
);
|
|
}
|
|
|
|
async removeGroupMember(groupId: string, endpointId: string): Promise<void> {
|
|
await this.request(
|
|
{ v: 1, type: "group.remove", rid: this.nextRid(), group_id: groupId, endpoint_id: endpointId },
|
|
false,
|
|
);
|
|
}
|
|
|
|
async leaveGroup(groupId: string): Promise<void> {
|
|
await this.request({ v: 1, type: "group.leave", rid: this.nextRid(), group_id: groupId }, false);
|
|
}
|
|
|
|
async transferGroup(groupId: string, endpointId: string): Promise<void> {
|
|
await this.request(
|
|
{ v: 1, type: "group.transfer", rid: this.nextRid(), group_id: groupId, endpoint_id: endpointId },
|
|
false,
|
|
);
|
|
}
|
|
|
|
async renameGroup(groupId: string, name: string): Promise<void> {
|
|
await this.request(
|
|
{ v: 1, type: "group.rename", rid: this.nextRid(), group_id: groupId, name },
|
|
false,
|
|
);
|
|
}
|
|
|
|
async dissolveGroup(groupId: string): Promise<void> {
|
|
await this.request({ v: 1, type: "group.dissolve", rid: this.nextRid(), group_id: groupId }, false);
|
|
}
|
|
|
|
async listGroups(cursor = "", limit = 0): Promise<unknown> {
|
|
const req: Record<string, unknown> = { v: 1, type: "group.list", rid: this.nextRid() };
|
|
if (cursor) req.cursor = cursor;
|
|
if (limit) req.limit = limit;
|
|
return this.request(req, false);
|
|
}
|
|
|
|
async getGroup(groupId: string, cursor = "", limit = 0): Promise<unknown> {
|
|
const req: Record<string, unknown> = {
|
|
v: 1,
|
|
type: "group.get",
|
|
rid: this.nextRid(),
|
|
group_id: groupId,
|
|
};
|
|
if (cursor) req.cursor = cursor;
|
|
if (limit) req.limit = limit;
|
|
return this.request(req, false);
|
|
}
|
|
|
|
async logout(): Promise<void> {
|
|
let err: unknown;
|
|
try {
|
|
await this.request({ v: 1, type: "self.logout", rid: this.nextRid() }, false);
|
|
} catch (e) {
|
|
err = e;
|
|
}
|
|
await this.teardown(false, new APIError("logged_out", "已退出"));
|
|
if (err) throw err;
|
|
}
|
|
|
|
async close(): Promise<void> {
|
|
await this.teardown(true, new APIError("closed", "已关闭"));
|
|
}
|
|
}
|
|
|
|
function bodyDecodedLen(b: Body): number {
|
|
if (b.enc === "base64") {
|
|
const bin = atob(b.data);
|
|
return bin.length;
|
|
}
|
|
return new TextEncoder().encode(b.data).length;
|
|
}
|
|
|
|
function utf8Len(s: string): number {
|
|
return new TextEncoder().encode(s).length;
|
|
}
|
|
|
|
function sleep(ms: number): Promise<void> {
|
|
return new Promise((r) => {
|
|
const t = setTimeout(r, ms);
|
|
t.unref?.();
|
|
});
|
|
}
|
|
|
|
export async function register(
|
|
connectOrRegisterURL: string,
|
|
registrationCode: string,
|
|
opt: RegisterOptions = {},
|
|
): Promise<RegisterResult> {
|
|
let regURL = connectOrRegisterURL;
|
|
try {
|
|
regURL = registerURLFromConnect(connectOrRegisterURL);
|
|
} catch {
|
|
/* 已是注册 URL */
|
|
}
|
|
const resp = await fetch(regURL, {
|
|
method: "POST",
|
|
headers: { "Content-Type": "application/json" },
|
|
body: marshalJSON({
|
|
registration_code: registrationCode,
|
|
id: opt.id ?? "",
|
|
login_password: opt.loginPassword ?? "",
|
|
name: opt.name ?? "",
|
|
talk_password: opt.talkPassword ?? "",
|
|
}),
|
|
});
|
|
const wrap = (await resp.json()) as {
|
|
ok: boolean;
|
|
data?: { id: string; login_password?: string };
|
|
error?: { code: string; message: string };
|
|
};
|
|
if (!wrap.ok) {
|
|
throw new APIError(wrap.error?.code ?? "bad_request", wrap.error?.message ?? "注册失败");
|
|
}
|
|
return { id: wrap.data!.id, loginPassword: wrap.data!.login_password };
|
|
}
|