Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
101 changes: 98 additions & 3 deletions apps/server/src/wsServer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import fs from "node:fs";
import os from "node:os";
import path from "node:path";

import * as NodeServices from "@effect/platform-node/NodeServices";
import { Effect, Exit, Layer, PubSub, Scope, Stream } from "effect";
import { describe, expect, it, afterEach, vi } from "vitest";
import { createServer } from "./wsServer";
Expand Down Expand Up @@ -38,7 +39,7 @@ import type {
TerminalWriteInput,
} from "@t3tools/contracts";
import { TerminalManager, type TerminalManagerShape } from "./terminal/Services/Manager";
import { SqlitePersistenceMemory } from "./persistence/Layers/Sqlite";
import { makeSqlitePersistenceLive, SqlitePersistenceMemory } from "./persistence/Layers/Sqlite";
import { SqlClient } from "effect/unstable/sql";
import { ProviderService, type ProviderServiceShape } from "./provider/Services/ProviderService";
import { Open, type OpenShape } from "./open";
Expand Down Expand Up @@ -448,23 +449,117 @@ describe("WebSocket Server", () => {

const ws = await connectWs(port);
connections.push(ws);
await waitForMessage(ws); // welcome
const welcome = (await waitForMessage(ws)) as WsPush; // welcome
expect(welcome.channel).toBe(WS_CHANNELS.serverWelcome);
expect(welcome.data).toEqual(
expect.objectContaining({
cwd: "/test/bootstrap-workspace",
projectName: "bootstrap-workspace",
bootstrapProjectId: expect.any(String),
bootstrapThreadId: expect.any(String),
}),
);

const snapshotResponse = await sendRequest(ws, ORCHESTRATION_WS_METHODS.getSnapshot);
expect(snapshotResponse.error).toBeUndefined();
const snapshot = snapshotResponse.result as {
projects: Array<{ workspaceRoot: string; title: string; defaultModel: string | null }>;
projects: Array<{
id: string;
workspaceRoot: string;
title: string;
defaultModel: string | null;
}>;
threads: Array<{
id: string;
projectId: string;
title: string;
model: string;
branch: string | null;
worktreePath: string | null;
}>;
};
const bootstrapProjectId = (welcome.data as { bootstrapProjectId?: string }).bootstrapProjectId;
const bootstrapThreadId = (welcome.data as { bootstrapThreadId?: string }).bootstrapThreadId;
expect(bootstrapProjectId).toBeDefined();
expect(bootstrapThreadId).toBeDefined();

expect(snapshot.projects).toEqual(
expect.arrayContaining([
expect.objectContaining({
id: bootstrapProjectId,
workspaceRoot: "/test/bootstrap-workspace",
title: "bootstrap-workspace",
defaultModel: "gpt-5-codex",
}),
]),
);
expect(snapshot.threads).toEqual(
expect.arrayContaining([
expect.objectContaining({
id: bootstrapThreadId,
projectId: bootstrapProjectId,
title: "New thread",
model: "gpt-5-codex",
branch: null,
worktreePath: null,
}),
]),
);
});

it("includes bootstrap ids in welcome when cwd project and thread already exist", async () => {
const stateDir = makeTempDir("t3code-state-bootstrap-existing-");
const persistenceLayer = makeSqlitePersistenceLive(
path.join(stateDir, "state.sqlite"),
).pipe(Layer.provide(NodeServices.layer)) as unknown as Layer.Layer<SqlClient.SqlClient, never>;
const cwd = "/test/bootstrap-existing";

server = await createTestServer({
cwd,
stateDir,
persistenceLayer,
autoBootstrapProjectFromCwd: true,
});
let addr = server.address();
let port = typeof addr === "object" && addr !== null ? addr.port : 0;
expect(port).toBeGreaterThan(0);

const firstWs = await connectWs(port);
connections.push(firstWs);
const firstWelcome = (await waitForMessage(firstWs)) as WsPush;
const firstBootstrapProjectId = (firstWelcome.data as { bootstrapProjectId?: string })
.bootstrapProjectId;
const firstBootstrapThreadId = (firstWelcome.data as { bootstrapThreadId?: string })
.bootstrapThreadId;
expect(firstBootstrapProjectId).toBeDefined();
expect(firstBootstrapThreadId).toBeDefined();

firstWs.close();
await closeTestServer();
server = null;

server = await createTestServer({
cwd,
stateDir,
persistenceLayer,
autoBootstrapProjectFromCwd: true,
});
addr = server.address();
port = typeof addr === "object" && addr !== null ? addr.port : 0;
expect(port).toBeGreaterThan(0);

const secondWs = await connectWs(port);
connections.push(secondWs);
const secondWelcome = (await waitForMessage(secondWs)) as WsPush;
expect(secondWelcome.channel).toBe(WS_CHANNELS.serverWelcome);
expect(secondWelcome.data).toEqual(
expect.objectContaining({
cwd,
projectName: "bootstrap-existing",
bootstrapProjectId: firstBootstrapProjectId,
bootstrapThreadId: firstBootstrapThreadId,
}),
);
});

it("logs outbound websocket push events in dev mode", async () => {
Expand Down
57 changes: 50 additions & 7 deletions apps/server/src/wsServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import {
ORCHESTRATION_WS_CHANNELS,
ORCHESTRATION_WS_METHODS,
ProjectId,
ThreadId,
TerminalEvent,
WS_CHANNELS,
WS_METHODS,
Expand Down Expand Up @@ -337,22 +338,59 @@ export const createServer = Effect.fn(function* (): Effect.fn.Return<

yield* Scope.provide(orchestrationReactor.start, subscriptionsScope);

let welcomeBootstrapProjectId: ProjectId | undefined;
let welcomeBootstrapThreadId: ThreadId | undefined;

if (autoBootstrapProjectFromCwd) {
yield* Effect.gen(function* () {
const snapshot = yield* projectionReadModelQuery.getSnapshot();
const existing = snapshot.projects.find((project) => project.workspaceRoot === cwd);
if (!existing) {
const existingProject = snapshot.projects.find(
(project) => project.workspaceRoot === cwd && project.deletedAt === null,
);
let bootstrapProjectId: ProjectId;
let bootstrapProjectDefaultModel: string;

if (!existingProject) {
const createdAt = new Date().toISOString();
const projectName = path.basename(cwd) || "project";
bootstrapProjectId = ProjectId.makeUnsafe(crypto.randomUUID());
const bootstrapProjectTitle = path.basename(cwd) || "project";
bootstrapProjectDefaultModel = "gpt-5-codex";
yield* orchestrationEngine.dispatch({
type: "project.create",
commandId: CommandId.makeUnsafe(crypto.randomUUID()),
projectId: ProjectId.makeUnsafe(crypto.randomUUID()),
title: projectName,
projectId: bootstrapProjectId,
title: bootstrapProjectTitle,
workspaceRoot: cwd,
defaultModel: "gpt-5-codex",
defaultModel: bootstrapProjectDefaultModel,
createdAt,
});
} else {
bootstrapProjectId = existingProject.id;
bootstrapProjectDefaultModel = existingProject.defaultModel ?? "gpt-5-codex";
}

const existingThread = snapshot.threads.find(
(thread) => thread.projectId === bootstrapProjectId && thread.deletedAt === null,
);
if (!existingThread) {
const createdAt = new Date().toISOString();
const threadId = ThreadId.makeUnsafe(crypto.randomUUID());
yield* orchestrationEngine.dispatch({
type: "thread.create",
commandId: CommandId.makeUnsafe(crypto.randomUUID()),
threadId,
projectId: bootstrapProjectId,
title: "New thread",
model: bootstrapProjectDefaultModel,
branch: null,
worktreePath: null,
createdAt,
});
welcomeBootstrapProjectId = bootstrapProjectId;
welcomeBootstrapThreadId = threadId;
} else {
welcomeBootstrapProjectId = bootstrapProjectId;
welcomeBootstrapThreadId = existingThread.id;
}
}).pipe(
Effect.mapError(
Expand Down Expand Up @@ -581,7 +619,12 @@ export const createServer = Effect.fn(function* (): Effect.fn.Return<
const welcome: WsPush = {
type: "push",
channel: WS_CHANNELS.serverWelcome,
data: { cwd, projectName },
data: {
cwd,
projectName,
...(welcomeBootstrapProjectId ? { bootstrapProjectId: welcomeBootstrapProjectId } : {}),
...(welcomeBootstrapThreadId ? { bootstrapThreadId: welcomeBootstrapThreadId } : {}),
},
};
logOutgoingPush(welcome, 1);
ws.send(JSON.stringify(welcome));
Expand Down
46 changes: 42 additions & 4 deletions apps/web/src/routes/__root.tsx
Original file line number Diff line number Diff line change
@@ -1,5 +1,11 @@
import { ThreadId } from "@t3tools/contracts";
import { Outlet, createRootRouteWithContext, type ErrorComponentProps } from "@tanstack/react-router";
import {
Outlet,
createRootRouteWithContext,
type ErrorComponentProps,
useNavigate,
useRouterState,
} from "@tanstack/react-router";
import { useEffect, useRef } from "react";
import { QueryClient, useQueryClient } from "@tanstack/react-query";

Expand Down Expand Up @@ -119,7 +125,12 @@ function errorDetails(error: unknown): string {
function EventRouter() {
const { dispatch } = useStore();
const queryClient = useQueryClient();
const navigate = useNavigate();
const pathname = useRouterState({ select: (state) => state.location.pathname });
const pathnameRef = useRef(pathname);
const lastConfigIssuesSignatureRef = useRef<string | null>(null);
const handledBootstrapThreadIdRef = useRef<string | null>(null);
pathnameRef.current = pathname;

useEffect(() => {
const api = readNativeApi();
Expand Down Expand Up @@ -180,8 +191,35 @@ function EventRouter() {
hasRunningSubprocess,
});
});
const unsubWelcome = onServerWelcome(() => {
void syncSnapshot();
const unsubWelcome = onServerWelcome((payload) => {
void (async () => {
await syncSnapshot();
if (disposed) {
return;
}

if (!payload.bootstrapProjectId || !payload.bootstrapThreadId) {
return;
}
dispatch({
type: "SET_PROJECT_EXPANDED",
projectId: payload.bootstrapProjectId,
expanded: true,
});

if (pathnameRef.current !== "/") {
return;
}
if (handledBootstrapThreadIdRef.current === payload.bootstrapThreadId) {
return;
}
await navigate({
to: "/$threadId",
params: { threadId: payload.bootstrapThreadId },
replace: true,
});
handledBootstrapThreadIdRef.current = payload.bootstrapThreadId;
})().catch(() => undefined);
Comment thread
cursor[bot] marked this conversation as resolved.
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const unsubServerConfigUpdated = onServerConfigUpdated((payload) => {
const signature = JSON.stringify(payload.issues);
Expand Down Expand Up @@ -232,7 +270,7 @@ function EventRouter() {
unsubWelcome();
unsubServerConfigUpdated();
};
}, [dispatch, queryClient]);
}, [dispatch, navigate, queryClient]);

return null;
}
Expand Down
9 changes: 9 additions & 0 deletions apps/web/src/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ type Action =
terminalId: string;
hasRunningSubprocess: boolean;
}
| { type: "SET_PROJECT_EXPANDED"; projectId: Project["id"]; expanded: boolean }
| { type: "TOGGLE_THREAD_TERMINAL"; threadId: ThreadId }
| { type: "SET_THREAD_TERMINAL_OPEN"; threadId: ThreadId; open: boolean }
| { type: "SET_THREAD_TERMINAL_HEIGHT"; threadId: ThreadId; height: number }
Expand Down Expand Up @@ -571,6 +572,14 @@ export function reducer(state: AppState, action: Action): AppState {
),
};

case "SET_PROJECT_EXPANDED":
return {
...state,
projects: state.projects.map((p) =>
p.id === action.projectId ? { ...p, expanded: action.expanded } : p,
),
};

case "TOGGLE_THREAD_TERMINAL":
return {
...state,
Expand Down
33 changes: 30 additions & 3 deletions apps/web/src/wsNativeApi.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,16 +65,41 @@ describe("wsNativeApi", () => {
emitPush(WS_CHANNELS.serverWelcome, payload);

expect(listener).toHaveBeenCalledTimes(1);
expect(listener).toHaveBeenCalledWith(payload);
expect(listener).toHaveBeenCalledWith(expect.objectContaining(payload));

const lateListener = vi.fn();
onServerWelcome(lateListener);

expect(lateListener).toHaveBeenCalledTimes(1);
expect(lateListener).toHaveBeenCalledWith(payload);
expect(lateListener).toHaveBeenCalledWith(expect.objectContaining(payload));
expect(warnSpy).not.toHaveBeenCalled();
});

it("preserves bootstrap ids from server.welcome payloads", async () => {
const { createWsNativeApi, onServerWelcome } = await import("./wsNativeApi");

createWsNativeApi();
const listener = vi.fn();
onServerWelcome(listener);

emitPush(WS_CHANNELS.serverWelcome, {
cwd: "/tmp/workspace",
projectName: "t3-code",
bootstrapProjectId: "project-1",
bootstrapThreadId: "thread-1",
});

expect(listener).toHaveBeenCalledTimes(1);
expect(listener).toHaveBeenCalledWith(
expect.objectContaining({
cwd: "/tmp/workspace",
projectName: "t3-code",
bootstrapProjectId: "project-1",
bootstrapThreadId: "thread-1",
}),
);
});

it("ignores invalid server.welcome payloads and keeps subscription active", async () => {
const warnSpy = vi.spyOn(console, "warn").mockImplementation(() => {});
const { createWsNativeApi, onServerWelcome } = await import("./wsNativeApi");
Expand All @@ -87,7 +112,9 @@ describe("wsNativeApi", () => {
emitPush(WS_CHANNELS.serverWelcome, { cwd: "/tmp/workspace", projectName: "t3-code" });

expect(listener).toHaveBeenCalledTimes(1);
expect(listener).toHaveBeenCalledWith({ cwd: "/tmp/workspace", projectName: "t3-code" });
expect(listener).toHaveBeenCalledWith(
expect.objectContaining({ cwd: "/tmp/workspace", projectName: "t3-code" }),
);
expect(warnSpy).toHaveBeenCalledTimes(1);
expect(warnSpy).toHaveBeenCalledWith("Dropped inbound WebSocket push payload", {
reason: "decode-failed",
Expand Down
4 changes: 3 additions & 1 deletion packages/contracts/src/ws.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { Schema, Struct } from "effect";
import { TrimmedNonEmptyString } from "./baseSchemas";
import { ProjectId, ThreadId, TrimmedNonEmptyString } from "./baseSchemas";

import {
ClientOrchestrationCommand,
Expand Down Expand Up @@ -163,5 +163,7 @@ export type WsResponse = typeof WsResponse.Type;
export const WsWelcomePayload = Schema.Struct({
cwd: TrimmedNonEmptyString,
projectName: TrimmedNonEmptyString,
bootstrapProjectId: Schema.optional(ProjectId),
bootstrapThreadId: Schema.optional(ThreadId),
});
export type WsWelcomePayload = typeof WsWelcomePayload.Type;