Skip to content

Commit dd48b03

Browse files
committed
feat(workspace): push worktree changes over websocket
1 parent d216f84 commit dd48b03

16 files changed

Lines changed: 794 additions & 11 deletions

File tree

apps/workspace/README.md

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,13 @@ set `VITE_UNIVER_LICENSE` at build time for any non-local deployment or to
125125
override it locally. Server, database, GitHub, and Discord settings are runtime
126126
values.
127127

128+
An authenticated Browser keeps one `/api/worktree-events` WebSocket open. AI or
129+
CLI Worktree writes publish a cache-invalidation signal only after the combined
130+
Collaboration and product operation completes, so active/processed task lists,
131+
details, sidebar counts, and Worktree-driven Node/Resource lists refresh without
132+
a page reload. The connection uses a one-time session ticket and carries no
133+
Worktree metadata or content.
134+
128135
## Docker
129136

130137
Build the image from the repository root:

apps/workspace/docs/application-design.md

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@ React
1919
│ ├── Permissions / Trash / Views
2020
│ ├── Worktrees
2121
│ └── Operations
22+
├── Worktree Change Feed
23+
│ └── Authenticated WebSocket cache invalidation
2224
2325
└── Univer Collaboration Client
2426
└── Collaboration Endpoint
@@ -30,6 +32,9 @@ React
3032
Product HTTP 与 Collaboration Endpoint 是两个入口,共享身份和权限模块,但不互相发起
3133
HTTP 回调。Application Module 通过进程内 Interface 调用 Collaboration Service。
3234

35+
Worktree Change Feed 是第三个窄入口,只传播产品缓存失效信号。它复用 Collaboration
36+
Endpoint 签发的一次性 Session Ticket,但不传播 snapshot、changeset、产品字段或权限结果。
37+
3338
## 数据库启动边界
3439

3540
`db/initialize.ts` 在创建业务 Repository 前打开数据库。磁盘数据库先经过可整体删除的
@@ -211,6 +216,15 @@ Collaboration Snapshot 重试、Node 元数据加载、Worktree Preview 和失
211216

212217
## Web 应用
213218

219+
AI 或 CLI 通过另一 Login Session 修改 Worktree 时,Worktree Module 在完整产品操作成功后
220+
向变更前后可发现该 Worktree 的在线用户发布 `worktreesChanged`。Browser 收到信号后使
221+
`["worktrees"]` Query 失效;连接建立时服务端先发送 `worktreeChangeFeedReady`,Browser
222+
同样执行一次失效,从而覆盖断线期间遗漏的 best-effort 通知。创建事件在产品 Worktree 与
223+
Operation 都已保存后发布;merge/discard 事件在 `processed_at` 与相关恢复状态收敛后发布。
224+
225+
该通道不替代 Collaboration Worktree 的 per-Worktree 状态连接,也不建立可供其他 Module
226+
任意发布的全局 Event Bus。
227+
214228
Web 应用的 Tree Row 总是 Node:
215229

216230
- `/nodes/{nodeId}` 是 Canonical URL;

apps/workspace/docs/architecture.md

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ SQLite 数据目录。
1414
| 服务端状态 | TanStack Query |
1515
| Web 本地状态 | React 本地状态 |
1616
| HTTP | Express 5 |
17+
| 产品实时失效 | 认证 WebSocket、TanStack Query invalidation |
1718
| 产品数据库 | Node `node:sqlite`、显式 SQL schema |
1819
| API 客户端 | `openapi-typescript``openapi-fetch` |
1920
| API 文档 | OpenAPI 3.1、Redocly CLI、Scalar |
@@ -129,6 +130,11 @@ shared → features → routes → app
129130
TanStack Query 管理 Session、Node、Resource、Recent、Trash、Permission、Worktree 和 Operation
130131
等服务端状态。Dialog、表单输入和当前选中项使用 React 本地状态,不引入额外全局状态库。
131132

133+
Worktree 页面使用一条用户级 `/api/worktree-events` WebSocket 接收粗粒度缓存失效信号。
134+
连接通过现有的一次性 Collaboration Session Ticket 认证;首次连接和收到变更信号时使
135+
`worktrees` 及可能被 Worktree 合入改变的 Node、Recent、Owned、Shared Query 失效并重新读取
136+
权威产品 API,不把 WebSocket payload 当作产品数据。
137+
132138
## 服务端
133139

134140
`main.ts` 读取配置、创建应用并监听端口;`app.ts` 创建 Express 实例并挂载中间件、业务
@@ -163,6 +169,11 @@ Univer 集中在 `integrations/univer`,向业务 Module 提供产品语义的
163169
跨产品数据库和 Collaboration Service 的写入由 `operations` Module 持久化和恢复,不用
164170
一次 SQLite transaction 假装覆盖两个系统。
165171

172+
Worktree Service 在 Collaboration 与产品写入均完成后调用专用 Change Feed。Change Feed
173+
不是通用应用 Event Bus;它只向该 Worktree 变更前后可发现的已连接用户发送不含 Worktree
174+
身份或内容的失效信号。实时发送失败不改变已经完成的产品写入,客户端重连后通过首帧统一
175+
失效查询,从产品 API 恢复当前状态。
176+
166177
## 产品数据库
167178

168179
产品数据库使用 Node `node:sqlite``db/schema.sql` 定义完整 V6 结构,`initialize.ts`

apps/workspace/server/src/app.ts

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import {
1818
createCollaborationGateway,
1919
type CollaborationGateway,
2020
} from "./integrations/univer/collaboration-gateway.js";
21+
import { createWorktreeChangeFeed } from "./integrations/realtime/worktree-change-feed.js";
2122
import {
2223
errorHandler,
2324
notFoundHandler,
@@ -209,6 +210,7 @@ export function createWorkspaceApplication(
209210
access,
210211
});
211212
const univerAssetsRepository = new UniverAssetsRepository(database);
213+
const worktreeChangeFeed = createWorktreeChangeFeed();
212214
const worktrees = createWorktreesModule({
213215
repository: new WorktreesRepository(database),
214216
access,
@@ -219,6 +221,7 @@ export function createWorkspaceApplication(
219221
: unavailableWorktreeBackend()),
220222
publishMergedAssets: (worktreeId, unitIds) =>
221223
univerAssetsRepository.publishWorktreeAssets(worktreeId, unitIds),
224+
onChanged: (change) => worktreeChangeFeed.publish(change),
222225
});
223226
const univerAssets = createUniverAssetsModule({
224227
repository: univerAssetsRepository,
@@ -245,6 +248,7 @@ export function createWorkspaceApplication(
245248
access,
246249
worktreeService: collaboration.worktreeService,
247250
worktrees,
251+
worktreeChangeFeed,
248252
})
249253
: null;
250254
invalidateRealtimeNodeAccess = () =>
@@ -346,11 +350,19 @@ export function createWorkspaceApplication(
346350
collaborationGateway?.attachWebSocket(server);
347351
},
348352
async closeRealtime() {
349-
await collaborationGateway?.dispose();
353+
try {
354+
await collaborationGateway?.dispose();
355+
} finally {
356+
await worktreeChangeFeed.dispose();
357+
}
350358
},
351359
async close() {
352360
try {
353-
await collaborationGateway?.dispose();
361+
try {
362+
await collaborationGateway?.dispose();
363+
} finally {
364+
await worktreeChangeFeed.dispose();
365+
}
354366
} finally {
355367
try {
356368
await exchange.dispose();
Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
import type { ISessionTicketStore } from "@univerjs-pro/collaboration-endpoint";
2+
import type {
3+
NodeTransportConnection,
4+
NodeTransportEndpoint,
5+
} from "@univerjs-pro/collaboration-transport-node";
6+
import type { WorktreeProductChange } from "../../modules/worktrees/index.js";
7+
8+
export const WORKTREE_CHANGE_FEED_PATH = "/api/worktree-events";
9+
10+
export interface WorktreeChangeFeed {
11+
endpoint(ticketStore: ISessionTicketStore): NodeTransportEndpoint;
12+
publish(change: WorktreeProductChange): void;
13+
dispose(): Promise<void>;
14+
}
15+
16+
export function createWorktreeChangeFeed(): WorktreeChangeFeed {
17+
const connections = new Map<
18+
string,
19+
Map<string, NodeTransportConnection>
20+
>();
21+
let disposed = false;
22+
let endpointCreated = false;
23+
24+
return {
25+
endpoint(ticketStore) {
26+
if (endpointCreated) {
27+
throw new Error("Worktree change feed endpoint is already created.");
28+
}
29+
endpointCreated = true;
30+
return {
31+
async handleUpgrade(context, next) {
32+
const url = new URL(
33+
context.incomingMessage.url ?? "/",
34+
"http://localhost"
35+
);
36+
if (url.pathname !== WORKTREE_CHANGE_FEED_PATH) {
37+
await next();
38+
return;
39+
}
40+
if (disposed) {
41+
context.reject(503, "Worktree change feed is unavailable");
42+
return;
43+
}
44+
const ticket = await ticketStore.consume(
45+
url.searchParams.get("sessionTicket") ?? ""
46+
);
47+
if (!ticket) {
48+
context.reject(401, "Invalid or expired session ticket");
49+
return;
50+
}
51+
context.accept({
52+
open({ connection }) {
53+
const userConnections =
54+
connections.get(ticket.userID) ?? new Map();
55+
userConnections.set(connection.id, connection);
56+
connections.set(ticket.userID, userConnections);
57+
safeSend(
58+
connection,
59+
JSON.stringify({ event: "worktreeChangeFeedReady" })
60+
);
61+
},
62+
message({ connection }) {
63+
connection.close(1003, "Worktree change feed is server-only");
64+
},
65+
close({ connection }) {
66+
removeConnection(ticket.userID, connection.id);
67+
},
68+
});
69+
},
70+
dispose: async () => {
71+
await dispose();
72+
},
73+
};
74+
},
75+
76+
publish(change) {
77+
if (disposed) return;
78+
const message = JSON.stringify({ event: "worktreesChanged" });
79+
for (const userId of change.audienceUserIds) {
80+
for (const connection of connections.get(userId)?.values() ?? []) {
81+
safeSend(connection, message);
82+
}
83+
}
84+
},
85+
86+
dispose,
87+
};
88+
89+
function removeConnection(userId: string, connectionId: string): void {
90+
const userConnections = connections.get(userId);
91+
userConnections?.delete(connectionId);
92+
if (userConnections?.size === 0) connections.delete(userId);
93+
}
94+
95+
async function dispose(): Promise<void> {
96+
if (disposed) return;
97+
disposed = true;
98+
for (const userConnections of connections.values()) {
99+
for (const connection of userConnections.values()) {
100+
connection.close(1001, "Worktree change feed disposed");
101+
}
102+
}
103+
connections.clear();
104+
}
105+
}
106+
107+
function safeSend(connection: NodeTransportConnection, message: string): void {
108+
void Promise.resolve(connection.send(message)).catch(() => {
109+
connection.close(1011, "Worktree change delivery failed");
110+
});
111+
}

apps/workspace/server/src/integrations/univer/collaboration-gateway.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,10 @@ import { json, Router, type RequestHandler } from "express";
3333
import type { AccessResolver, ResourceAccess, UnitType } from "../../modules/access/index.js";
3434
import type { IdentityModule } from "../../modules/identity/index.js";
3535
import type { WorktreesModule } from "../../modules/worktrees/index.js";
36+
import {
37+
WORKTREE_CHANGE_FEED_PATH,
38+
type WorktreeChangeFeed,
39+
} from "../realtime/worktree-change-feed.js";
3640

3741
const OK_ERROR = { code: ErrorCode.OK, message: "" };
3842

@@ -49,13 +53,15 @@ export function createCollaborationGateway(options: {
4953
readonly access: AccessResolver;
5054
readonly worktreeService: UniverCollabWorktreeService;
5155
readonly worktrees: WorktreesModule;
56+
readonly worktreeChangeFeed: WorktreeChangeFeed;
5257
}): CollaborationGateway {
5358
const {
5459
service,
5560
identity,
5661
access,
5762
worktreeService,
5863
worktrees,
64+
worktreeChangeFeed,
5965
} = options;
6066
const ticketStore = new MemorySessionTicketStore();
6167
const endpoint = new UniverCollabEndpoint(service, { ticketStore });
@@ -222,6 +228,7 @@ export function createCollaborationGateway(options: {
222228
await next();
223229
});
224230
transport.register(worktreeEndpoint);
231+
transport.register(worktreeChangeFeed.endpoint(ticketStore));
225232
transport.register(trackConnections(endpoint, nodeAccessConnections));
226233

227234
const router = Router();
@@ -271,7 +278,8 @@ export function createCollaborationGateway(options: {
271278
const url = new URL(request.url ?? "/", "http://localhost");
272279
if (
273280
url.pathname !== "/universer-api/comb/connect" &&
274-
!url.pathname.startsWith("/universer-api/worktrees/")
281+
!url.pathname.startsWith("/universer-api/worktrees/") &&
282+
url.pathname !== WORKTREE_CHANGE_FEED_PATH
275283
) {
276284
return;
277285
}

apps/workspace/server/src/modules/worktrees/index.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ export type {
1010
WorktreeDetail,
1111
WorktreeKind,
1212
WorktreeOperationView,
13+
WorktreeProductChange,
1314
WorktreesModule,
1415
WorktreeState,
1516
WorktreeSummary,

apps/workspace/server/src/modules/worktrees/worktrees.repository.ts

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,42 @@ export interface WorktreeOperationRow {
7373
export class WorktreesRepository {
7474
constructor(private readonly _database: WorkspaceDatabase) {}
7575

76+
audience(worktreeId: string): string[] {
77+
const rows = this._database.connection
78+
.prepare(
79+
`SELECT audience.user_id
80+
FROM (
81+
SELECT worktree.creator_user_id AS user_id
82+
FROM worktrees AS worktree
83+
WHERE worktree.id = ?
84+
85+
UNION
86+
87+
SELECT team.owner_user_id AS user_id
88+
FROM worktrees AS worktree
89+
JOIN spaces AS team ON team.id = worktree.team_space_id
90+
WHERE worktree.id = ?
91+
92+
UNION
93+
94+
SELECT membership.user_id
95+
FROM worktrees AS worktree
96+
JOIN space_members AS membership
97+
ON membership.space_id = worktree.team_space_id
98+
WHERE worktree.id = ?
99+
AND (
100+
worktree.visibility = 'space'
101+
OR membership.role = 'admin'
102+
)
103+
) AS audience
104+
ORDER BY audience.user_id`
105+
)
106+
.all(worktreeId, worktreeId, worktreeId) as unknown as Array<{
107+
readonly user_id: string;
108+
}>;
109+
return rows.map((row) => row.user_id);
110+
}
111+
76112
listCandidates(input: {
77113
readonly userId: string;
78114
readonly scope: "active" | "processed";

0 commit comments

Comments
 (0)