Skip to content

Commit 0d2c76a

Browse files
committed
Stabilize scoped routing and reconnect verification
1 parent 92e373b commit 0d2c76a

5 files changed

Lines changed: 72 additions & 15 deletions

File tree

‎src/server/index.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1816,7 +1816,10 @@ async function bootstrap(): Promise<void> {
18161816
registerWorkflowDispatcher((bot, channelId, triggerId, threadRootId) => runBot(bot, channelId, triggerId, threadRootId, true));
18171817
reactivateComputersAfterPreparedRemoval();
18181818
prepareMnemosyneRuntime();
1819-
await startRoutingEngine((activity) => broadcastAdmins({ type: "routing_activity", activity }));
1819+
await startRoutingEngine((activity, ownerUserId) => {
1820+
if (ownerUserId) sendToUsers([ownerUserId], { type: "routing_activity", activity });
1821+
else broadcastAdmins({ type: "routing_activity", activity });
1822+
});
18201823
ensureImageGenerationSkill();
18211824
await internalRoutingProviderId();
18221825
resumeQueuedAgentTurns();

‎src/server/routing.ts‎

Lines changed: 30 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -118,6 +118,8 @@ type UserGatewayRuntime = {
118118
port: number;
119119
gateway: UserGateway;
120120
router: RoutingRuntime["router"];
121+
requestActivity: RoutingRuntime["requestActivity"];
122+
activityUnsubscribe: () => void;
121123
};
122124

123125
// ReRouted's generic OpenAI-compatible adapter historically assumed every
@@ -151,8 +153,9 @@ const INTERNAL_GATEWAY_KEY_SCOPE = "1helm-internal";
151153
let runtime: RoutingRuntime | null = null;
152154
let starting: Promise<RoutingRuntime> | null = null;
153155
let activityUnsubscribe: (() => void) | null = null;
154-
let onActivity: ((activity?: unknown) => void) | null = null;
156+
let onActivity: ((activity?: unknown, userId?: number) => void) | null = null;
155157
const recentActivity: unknown[] = [];
158+
const recentUserActivity = new Map<number, unknown[]>();
156159
type OauthCompletion = { connected: boolean; account?: Record<string, unknown>; error?: string };
157160
const oauthWatchers = new Map<string, NodeJS.Timeout>();
158161
const oauthCompletions = new Map<string, OauthCompletion>();
@@ -161,6 +164,15 @@ const oauthOwners = new Map<string, { userId: number; type: string; providerIds:
161164
const activeOauthByFamily = new Map<string, string>();
162165
const userGateways = new Map<number, UserGatewayRuntime>();
163166

167+
function publishRoutingActivity(activity: unknown, userId = 0): void {
168+
const history = userId > 0
169+
? (recentUserActivity.get(userId) || (recentUserActivity.set(userId, []), recentUserActivity.get(userId)!))
170+
: recentActivity;
171+
history.unshift(activity);
172+
if (history.length > 30) history.length = 30;
173+
onActivity?.(activity, userId);
174+
}
175+
164176
const oauthKey = (userId: number, type: string): string => `${userId}:${type}`;
165177
const oauthFamily = (type: string): string => ["chatgpt", "codex"].includes(type) ? "chatgpt" : type;
166178

@@ -276,10 +288,17 @@ async function ensureUserGateway(userId: number): Promise<UserGatewayRuntime> {
276288
},
277289
} as RoutingRuntime["router"];
278290
const gateway = createGateway({ store, router, requestActivity });
291+
const activityUnsubscribe = requestActivity.subscribe((activity) => publishRoutingActivity(activity, userId));
279292
const desiredPort = await choosePort(endpoint.port, "127.0.0.1", false);
280-
const address = await gateway.start(desiredPort, "127.0.0.1");
293+
let address: { port: number; host: string };
294+
try {
295+
address = await gateway.start(desiredPort, "127.0.0.1");
296+
} catch (error) {
297+
activityUnsubscribe();
298+
throw error;
299+
}
281300
if (address.port !== endpoint.port) run("UPDATE user_routing_endpoints SET port=?,updated=? WHERE user_id=?", address.port, now(), userId);
282-
const created = { userId, port: address.port, gateway, router };
301+
const created = { userId, port: address.port, gateway, router, requestActivity, activityUnsubscribe };
283302
userGateways.set(userId, created);
284303
return created;
285304
}
@@ -701,7 +720,7 @@ async function choosePort(preferred: number, host: string, honorConfiguredPort =
701720
});
702721
}
703722

704-
export async function startRoutingEngine(activityCallback?: (activity?: unknown) => void): Promise<RoutingRuntime> {
723+
export async function startRoutingEngine(activityCallback?: (activity?: unknown, userId?: number) => void): Promise<RoutingRuntime> {
705724
if (runtime) return runtime;
706725
if (starting) return starting;
707726
onActivity = activityCallback || null;
@@ -723,11 +742,7 @@ export async function startRoutingEngine(activityCallback?: (activity?: unknown)
723742
if (port !== Number(config.port || 4949)) target.store.update((current) => { current.port = port; });
724743
await target.start({ port, host });
725744
ensureInternalProvider(target);
726-
activityUnsubscribe = target.requestActivity.subscribe((activity) => {
727-
recentActivity.unshift(activity);
728-
if (recentActivity.length > 30) recentActivity.length = 30;
729-
onActivity?.(activity);
730-
});
745+
activityUnsubscribe = target.requestActivity.subscribe((activity) => publishRoutingActivity(activity));
731746
runtime = target;
732747
return target;
733748
})().finally(() => { starting = null; });
@@ -746,7 +761,10 @@ export async function stopRoutingEngine(): Promise<void> {
746761
activeOauthByFamily.clear();
747762
const gateways = [...userGateways.values()];
748763
userGateways.clear();
764+
for (const entry of gateways) entry.activityUnsubscribe();
749765
await Promise.all(gateways.map((entry) => entry.gateway.stop().catch(() => undefined)));
766+
recentActivity.length = 0;
767+
recentUserActivity.clear();
750768
if (target) await target.close({ drainMs: 10_000 });
751769
}
752770

@@ -955,7 +973,9 @@ export async function routingState(userId = 0, isAdmin = true): Promise<Record<s
955973
visibility: visibilityOf(combo),
956974
mine: Number(combo.ownerUserId || 0) === userId,
957975
}));
958-
return { ...state, providers, combos, activeRequests: runtime?.requestActivity.snapshot() || [], recentActivity, imageGenerationEnabled: enabled, apiKey: state.apiKey ? "" : state.apiKey, apiKeys: undefined, scope: isAdmin ? "captain" : "member" };
976+
const personalActivity = userId ? recentUserActivity.get(userId) || [] : recentActivity;
977+
const activeRequests = userId ? userGateways.get(userId)?.requestActivity.snapshot() || [] : runtime?.requestActivity.snapshot() || [];
978+
return { ...state, providers, combos, activeRequests, recentActivity: personalActivity, imageGenerationEnabled: enabled, apiKey: state.apiKey ? "" : state.apiKey, apiKeys: undefined, scope: isAdmin ? "captain" : "member" };
959979
}
960980

961981
export async function routingCredentials(userId = 0, isAdmin = true): Promise<Record<string, unknown>> {

‎test/connectors.mjs‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,9 +35,15 @@ test("stopping a connector cancels automatic relaunch while preserving its crede
3535
const beforeStop = (await readFile(log, "utf8")).trim().split("\n").length;
3636
assert.ok(beforeStop >= 2, "a failed desired connector relaunches");
3737
connectors.stopTunnelConnector("workspace");
38+
// A process already returned by spawn() can reach its first instruction
39+
// while stop is delivering SIGTERM. Establish the stopped boundary after
40+
// that in-flight child settles, then prove no timer can launch another one.
41+
await new Promise((resolve) => setTimeout(resolve, 75));
42+
const stoppedBoundary = (await readFile(log, "utf8")).trim().split("\n").length;
3843
await new Promise((resolve) => setTimeout(resolve, 140));
3944
const afterStop = (await readFile(log, "utf8")).trim().split("\n").length;
40-
assert.equal(afterStop, beforeStop, "disabled connector does not relaunch");
45+
assert.ok(stoppedBoundary >= beforeStop, "stop accounts for any child that was already spawned");
46+
assert.equal(afterStop, stoppedBoundary, "disabled connector does not relaunch after the stopped boundary");
4147
assert.ok(connectors.connectorCredential("workspace"), "disabling leaves the reserved connector credential available for re-enable");
4248

4349
process.env.HELM_CONNECTOR_TEST_HOLD = "1";

‎test/routing.mjs‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,6 +162,14 @@ test("embedded provider fabric powers 1Helm agents and its public endpoint", { t
162162
assert.notEqual(guestCredentials.personalPort, captainCredentialsBefore.personalPort, "each signed-in member receives a dedicated host-side endpoint port");
163163
assert.equal((await json(`http://127.0.0.1:${appPort}/v1/models`, guestExternal.key.key)).data.some((model) => model.id.includes("mock-large")), true, "member key routes through personal plus shared providers");
164164
await json(`http://127.0.0.1:${appPort}/api/routing/action`, guestToken, { method: "POST", body: JSON.stringify({ action: "app:set-provider-visibility", payload: { id: guestKeyed.id, visibility: "personal" } }) });
165+
const guestPrivateModel = (await json(`http://127.0.0.1:${appPort}/api/routing/models`, guestToken)).models.find((model) => model.providerName === "Crew private").id;
166+
await json(`http://127.0.0.1:${appPort}/v1/chat/completions`, guestExternal.key.key, {
167+
method: "POST", body: JSON.stringify({ model: guestPrivateModel, stream: false, messages: [{ role: "user", content: "Private member activity" }] }),
168+
});
169+
const guestActivity = (await json(`http://127.0.0.1:${appPort}/api/routing/state`, guestToken)).recentActivity;
170+
const captainActivity = (await json(`http://127.0.0.1:${appPort}/api/routing/state`, token)).recentActivity;
171+
assert.equal(guestActivity.some((event) => event.request?.providerName === "Crew private"), true, "personal endpoint activity is visible to its owner");
172+
assert.equal(captainActivity.some((event) => event.request?.providerName === "Crew private"), false, "personal endpoint activity is not disclosed to the Captain or another member");
165173
assert.equal((await json(`http://127.0.0.1:${appPort}/api/routing/state`, token)).combos.some((entry) => entry.name === "crew-shared-route"), false, "a shared route stops resolving for teammates when its provider becomes private");
166174
const guestOwnedSharedRoute = (await json(`http://127.0.0.1:${appPort}/api/routing/state`, guestToken)).combos.find((entry) => entry.name === "crew-shared-route");
167175
await json(`http://127.0.0.1:${appPort}/api/workspace/model-policy`, guestToken, { method: "PATCH", body: JSON.stringify({ model: "crew-shared-route", personal: true }) });

‎test/terminal-reconnect-browser.mjs‎

Lines changed: 23 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,30 @@
11
import assert from "node:assert/strict";
22
import { spawn } from "node:child_process";
33
import { createServer } from "node:net";
4-
import { mkdtempSync, rmSync } from "node:fs";
4+
import { accessSync, constants as fsConstants, mkdtempSync, rmSync } from "node:fs";
55
import { tmpdir } from "node:os";
66
import { join } from "node:path";
77
import test from "node:test";
88
import puppeteer from "puppeteer";
99

1010
const root = new URL("..", import.meta.url).pathname;
11+
function browserExecutable() {
12+
const configured = process.env.PUPPETEER_EXECUTABLE_PATH;
13+
if (configured) {
14+
try { accessSync(configured, fsConstants.X_OK); return configured; } catch { /* use discovery */ }
15+
}
16+
try {
17+
const bundled = puppeteer.executablePath();
18+
accessSync(bundled, fsConstants.X_OK);
19+
return bundled;
20+
} catch { /* no bundled browser */ }
21+
for (const candidate of ["/usr/bin/google-chrome", "/usr/bin/google-chrome-stable", "/usr/bin/chromium", "/usr/bin/chromium-browser"]) {
22+
try { accessSync(candidate, fsConstants.X_OK); return candidate; } catch { /* try next */ }
23+
}
24+
return null;
25+
}
26+
27+
const executablePath = browserExecutable();
1128
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
1229
const freePort = () => new Promise((resolve, reject) => {
1330
const server = createServer();
@@ -28,7 +45,10 @@ const waitFor = async (fn, label, timeout = 15_000) => {
2845
throw last instanceof Error ? last : new Error(`Timed out waiting for ${label}`);
2946
};
3047

31-
test("browser silently reconnects a dropped terminal socket to the same live shell", { timeout: 45_000 }, async () => {
48+
test("browser silently reconnects a dropped terminal socket to the same live shell", {
49+
timeout: 45_000,
50+
skip: executablePath ? false : "No local Chrome executable; the mandatory transport/keepalive contract still runs.",
51+
}, async () => {
3252
const dataDir = mkdtempSync(join(tmpdir(), "1helm-terminal-reconnect-"));
3353
const appPort = await freePort();
3454
const mockPort = await freePort();
@@ -62,7 +82,7 @@ test("browser silently reconnects a dropped terminal socket to the same live she
6282
await api("/api/setup/complete", registration.token, { name: "Reconnect Test", terminals_enabled: true, provider_id: provider.provider.id, model: "mock-large" });
6383
const channel = (await api("/api/channels", registration.token, { name: "terminal-reconnect", purpose: "Prove terminal session continuity." })).channel;
6484

65-
browser = await puppeteer.launch({ headless: true, args: ["--no-sandbox", "--disable-setuid-sandbox"] });
85+
browser = await puppeteer.launch({ executablePath, headless: true, args: ["--no-sandbox", "--disable-setuid-sandbox"] });
6686
const page = await browser.newPage();
6787
await page.evaluateOnNewDocument(() => {
6888
const NativeWebSocket = window.WebSocket;

0 commit comments

Comments
 (0)