fix(coding-agent/tools): serialized tab grouping operations in the relay bridge
- Serialized group and ungroup operations to prevent duplicate tab group creation races. - Queued and serially drained tab grouping requests in the relay bridge to prevent overlapping RPCs. - Mirrored tab group titles to session storage and healed duplicate groups during background service worker recovery. - Renamed run-cancellation utility to run-scope and updated corresponding module and test references.
This commit is contained in:
@@ -4,4 +4,8 @@
|
||||
|
||||
### Added
|
||||
|
||||
- Initial release: Chrome MV3 extension that lets the omp browser tool attach to and drive the user's existing tabs through `chrome.debugger`. The companion CDP relay lives in the omp CLI (`omp browser-relay`); this package builds the extension zip for GitHub releases and generates the embedded install assets consumed by `omp browser-relay install`. Controllable tabs are gathered into a per-window "omp" tab group while the relay is connected.
|
||||
- Initial release: Chrome MV3 extension that lets the omp browser tool attach to and drive the user's existing tabs through `chrome.debugger`. The companion CDP relay lives in the omp CLI (`omp browser-relay`); this package builds the extension zip for GitHub releases and generates the embedded install assets consumed by `omp browser-relay install`. Tabs the agent actively drives are gathered into a per-window "omp" tab group while the relay is connected.
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed duplicate "omp" tab groups: `group`/`ungroup` RPCs are now serialized in the service worker (Chrome's query→create→set-title sequence is not atomic), and `groupTabs` folds stray same-title groups left by earlier races back into the canonical per-window group. The omp group title is also mirrored to `chrome.storage.session`, so a service worker restarted between grouping and disconnect can still dissolve the group instead of leaving it behind.
|
||||
|
||||
@@ -11,7 +11,7 @@ The companion relay server lives in the omp CLI (`omp browser-relay`, see `packa
|
||||
|
||||
That's it: the relay server auto-starts under omp's profile-independent global daemon broker the first time the browser tool needs it. Every relay consumer holds a broker lease, so one project exiting cannot interrupt another; the server stops after the last consumer across all projects exits. The extension badge turns **on** when connected. Run `omp browser-relay` manually only for `--token`, `--no-group`, or a non-default port — a relay already serving the port is adopted, never fought over.
|
||||
|
||||
`app.target` picks a specific tab by URL/title substring; without it, omp adopts the visible tab without stealing focus. While connected, every tab omp can control is gathered into a per-window **"omp" tab group** (cyan) — dissolved on disconnect; pinned tabs, tabs in your own groups, and tabs you drag out are left alone. Disable with `omp browser-relay --no-group`.
|
||||
`app.target` picks a specific tab by URL/title substring; without it, omp adopts the visible tab without stealing focus. Tabs omp is **actively driving** are gathered into a per-window **"omp" tab group** (cyan) — released when omp lets go of the tab and dissolved on disconnect; the rest of your tabs, pinned tabs, tabs in your own groups, and tabs you drag out are left alone. Disable with `omp browser-relay --no-group`.
|
||||
|
||||
## Development
|
||||
|
||||
|
||||
@@ -47,12 +47,25 @@ function snapshot(tab: ChromeTab): TabSnapshot | null {
|
||||
};
|
||||
}
|
||||
|
||||
/** Title of the omp tab group, remembered so a relay disconnect can dissolve it. */
|
||||
/** Title of the omp tab group; mirrored to session storage so a restarted service worker can still dissolve it. */
|
||||
let ompGroupTitle: string | null = null;
|
||||
|
||||
/**
|
||||
* Serialize group mutations. Chrome's query→group→set-title sequence is not
|
||||
* atomic: two concurrent runs both miss the not-yet-titled group and mint
|
||||
* duplicate "omp" groups in the same window.
|
||||
*/
|
||||
let groupOps: Promise<unknown> = Promise.resolve();
|
||||
function enqueueGroupOp<T>(fn: () => Promise<T>): Promise<T> {
|
||||
const result = groupOps.then(fn, fn);
|
||||
groupOps = result.catch(() => {});
|
||||
return result;
|
||||
}
|
||||
|
||||
/** Move tabs into the per-window omp group, creating or reusing it by title. */
|
||||
async function groupTabs(tabIds: number[], title: string, color: string): Promise<{ grouped: Record<string, number> }> {
|
||||
ompGroupTitle = title;
|
||||
void chrome.storage.session.set({ ompGroupTitle: title });
|
||||
const byWindow = new Map<number, number[]>();
|
||||
for (const tabId of tabIds) {
|
||||
try {
|
||||
@@ -69,7 +82,19 @@ async function groupTabs(tabIds: number[], title: string, color: string): Promis
|
||||
const grouped: Record<string, number> = {};
|
||||
for (const [windowId, ids] of byWindow) {
|
||||
const existing = await chrome.tabGroups.query({ title, windowId });
|
||||
const groupId = await chrome.tabs.group(existing[0] ? { tabIds: ids, groupId: existing[0].id } : { tabIds: ids });
|
||||
let groupId: number;
|
||||
if (existing[0]) {
|
||||
groupId = existing[0].id;
|
||||
// Heal duplicate same-title groups left behind by older races.
|
||||
for (const dupe of existing.slice(1)) {
|
||||
const dupeTabs = await chrome.tabs.query({ groupId: dupe.id });
|
||||
const dupeIds = dupeTabs.map(tab => tab.id).filter(id => id !== undefined);
|
||||
if (dupeIds.length > 0) await chrome.tabs.group({ tabIds: dupeIds, groupId });
|
||||
}
|
||||
await chrome.tabs.group({ tabIds: ids, groupId });
|
||||
} else {
|
||||
groupId = await chrome.tabs.group({ tabIds: ids });
|
||||
}
|
||||
await chrome.tabGroups.update(groupId, { title, color });
|
||||
for (const id of ids) grouped[String(id)] = groupId;
|
||||
}
|
||||
@@ -78,6 +103,11 @@ async function groupTabs(tabIds: number[], title: string, color: string): Promis
|
||||
|
||||
/** Dissolve every omp-titled group (relay disconnected or asked us to release tabs). */
|
||||
async function restoreGroups(): Promise<void> {
|
||||
if (!ompGroupTitle) {
|
||||
// Service worker restarted since the last group op; recover the title.
|
||||
const stored = await chrome.storage.session.get({ ompGroupTitle: "" }).catch(() => ({ ompGroupTitle: "" }));
|
||||
ompGroupTitle = typeof stored.ompGroupTitle === "string" && stored.ompGroupTitle ? stored.ompGroupTitle : null;
|
||||
}
|
||||
if (!ompGroupTitle) return;
|
||||
const groups = await chrome.tabGroups.query({ title: ompGroupTitle }).catch(() => []);
|
||||
for (const group of groups) {
|
||||
@@ -151,9 +181,9 @@ async function runRpc(msg: Extract<RelayToExtMessage, { t: "rpc" }>): Promise<un
|
||||
return {};
|
||||
}
|
||||
case "group":
|
||||
return await groupTabs(msg.tabIds, msg.title, msg.color);
|
||||
return await enqueueGroupOp(() => groupTabs(msg.tabIds, msg.title, msg.color));
|
||||
case "ungroup":
|
||||
await chrome.tabs.ungroup(msg.tabIds).catch(() => {});
|
||||
await enqueueGroupOp(() => chrome.tabs.ungroup(msg.tabIds).catch(() => {}));
|
||||
return {};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -81,6 +81,10 @@ declare const chrome: {
|
||||
get(keys: Record<string, unknown>): Promise<Record<string, unknown>>;
|
||||
set(items: Record<string, unknown>): Promise<void>;
|
||||
};
|
||||
session: {
|
||||
get(keys: Record<string, unknown>): Promise<Record<string, unknown>>;
|
||||
set(items: Record<string, unknown>): Promise<void>;
|
||||
};
|
||||
onChanged: ChromeEvent<(changes: Record<string, unknown>, areaName: string) => void>;
|
||||
};
|
||||
alarms: {
|
||||
|
||||
@@ -4,13 +4,17 @@
|
||||
|
||||
### Added
|
||||
|
||||
- Added a `relay` browser mode that drives the user's own Chrome tabs through a local CDP relay plus the OMP Browser Relay extension: `omp browser-relay install` writes the bundled extension to disk, and `browser.relay` / `browser.relayUrl` (or per-call `app.relay`) route the browser tool through it. The relay server auto-starts under a profile-independent global daemon broker when the browser tool needs it; every relay consumer holds a broker lease, so the fixed-port singleton stops only after its last consumer across all projects exits. `omp browser-relay` remains available for `--token`/`--no-group`/custom ports, and a relay already serving the port is adopted. It multiplexes the supervisor and per-tab worker puppeteer connections over the single `chrome.debugger` attachment Chrome allows per tab, and gathers controllable tabs into a per-window "omp" tab group (dissolved on disconnect, never re-grouping tabs the user pulls out).
|
||||
- Added a `relay` browser mode that drives the user's own Chrome tabs through a local CDP relay plus the OMP Browser Relay extension: `omp browser-relay install` writes the bundled extension to disk, and `browser.relay` / `browser.relayUrl` (or per-call `app.relay`) route the browser tool through it. The relay server auto-starts under a profile-independent global daemon broker when the browser tool needs it; every relay consumer holds a broker lease, so the fixed-port singleton stops only after its last consumer across all projects exits. `omp browser-relay` remains available for `--token`/`--no-group`/custom ports, and a relay already serving the port is adopted. It multiplexes the supervisor and per-tab worker puppeteer connections over the single `chrome.debugger` attachment Chrome allows per tab, and gathers only the tabs the agent actively drives into a per-window "omp" tab group (released when the last client lets go of the tab, dissolved on disconnect, never re-grouping tabs the user pulls out).
|
||||
- Added individual-window computer use: a capture-free discovery call lists current window ids alongside a `desktop` target, then the model can capture and control one selected window without stealing focus or moving the user's pointer.
|
||||
|
||||
### Changed
|
||||
|
||||
- Exposed `computer` through its window-aware function schema for every model, including models with provider-native Computer Use support, because native computer declarations cannot carry the required window target.
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed the browser relay creating duplicate "omp" tab groups: the bridge now keeps at most one group RPC in flight (a queued drain replaces fire-and-forget per-tab requests), so concurrent requests can no longer race the extension's non-atomic query→create→set-title sequence in the same window. Also fixed an extension reconnect (relay daemon restart, service-worker recycle) being misread as the user dragging every tab out of the omp group — grouping state is reset when the extension socket closes, so tabs regroup on the next hello instead of being permanently opted out.
|
||||
|
||||
## [17.2.4] - 2026-08-01
|
||||
|
||||
### Added
|
||||
|
||||
@@ -22,7 +22,7 @@ export interface BrowserRelayCommandArgs {
|
||||
token?: string;
|
||||
/** Install target directory; defaults to ~/.omp/browser-relay/extension. */
|
||||
dir?: string;
|
||||
/** Gather controllable tabs into an 'omp' Chrome tab group (default true). */
|
||||
/** Gather tabs the agent actively drives into an 'omp' Chrome tab group (default true). */
|
||||
group?: boolean;
|
||||
verbose?: boolean;
|
||||
}
|
||||
|
||||
@@ -13,11 +13,11 @@ import { type AriaSnapshotOptions, assertSelectorString, buildAriaSnapshotScript
|
||||
import { DEFAULT_VIEWPORT } from "../launch";
|
||||
import { extractReadableFromHtml, type ReadableFormat } from "../readable";
|
||||
import {
|
||||
bindBrowserRunFacade,
|
||||
bindRunFacade,
|
||||
resolvePredicateTimeout,
|
||||
type WaitPredicateOptions,
|
||||
waitForBrowserRun,
|
||||
} from "../run-cancellation";
|
||||
waitForRun,
|
||||
} from "../../run-scope";
|
||||
import { cloneSafe, RunOutput } from "../run-output";
|
||||
import type { Observation, ReadyInfo, RunResultOk, ScreenshotResult, SessionSnapshot } from "../tab-protocol";
|
||||
import {
|
||||
@@ -1345,16 +1345,16 @@ export async function runCmuxCode(tab: CmuxTab, opts: RunCmuxCodeOptions): Promi
|
||||
// Keep both inside try so a concurrent in-process eval/browser run surfaces as
|
||||
// a rejected promise the supervisor can report, never an unhandled rejection.
|
||||
runtime.setCwd(opts.snapshot.cwd);
|
||||
const runTab = bindBrowserRunFacade(tab, signal);
|
||||
const runTab = bindRunFacade(tab, signal);
|
||||
runtime.setRunScope({
|
||||
page: bindBrowserRunFacade(tab.page, signal),
|
||||
browser: bindBrowserRunFacade(tab.browser, signal),
|
||||
page: bindRunFacade(tab.page, signal),
|
||||
browser: bindRunFacade(tab.browser, signal),
|
||||
tab: runTab,
|
||||
assert: (cond: unknown, text?: string): void => {
|
||||
if (!cond) throw new ToolError(text ?? "Assertion failed");
|
||||
},
|
||||
wait: (msOrPredicate: number | (() => unknown), waitOpts?: WaitPredicateOptions): Promise<unknown> =>
|
||||
waitForBrowserRun(
|
||||
waitForRun(
|
||||
msOrPredicate,
|
||||
signal,
|
||||
typeof msOrPredicate === "number"
|
||||
|
||||
@@ -55,6 +55,8 @@ class CdpConnection {
|
||||
autoAttach = false;
|
||||
/** Minted pseudo-sessions owned by this connection. */
|
||||
readonly sessions = new Map<string, SessionRef>();
|
||||
/** Tabs this connection claimed as drive targets (`OMP.claimTarget` / `Target.createTarget`). */
|
||||
readonly claims = new Set<number>();
|
||||
|
||||
constructor(
|
||||
readonly id: number,
|
||||
@@ -159,13 +161,17 @@ export class RelayBridge {
|
||||
/** Real child session id → owning tab, learned from `Target.attachedToTarget` events. */
|
||||
#realSessionTabs = new Map<string, number>();
|
||||
#log: (message: string, data?: Record<string, unknown>) => void;
|
||||
/** Tab-group appearance for controllable tabs; null disables grouping. */
|
||||
/** Tab-group appearance for driven tabs; null disables grouping. */
|
||||
#group: { title: string; color: string } | null;
|
||||
/** Tabs awaiting the next group RPC; drained one batch at a time. */
|
||||
#groupQueue: TabState[] = [];
|
||||
/** True while {@link #drainGroupQueue} runs — group RPCs must never overlap. */
|
||||
#groupDraining = false;
|
||||
|
||||
constructor(
|
||||
opts: {
|
||||
log?: (message: string, data?: Record<string, unknown>) => void;
|
||||
/** Group controllable tabs under one per-window Chrome tab group. */
|
||||
/** Group tabs the agent actively drives under one per-window Chrome tab group. */
|
||||
group?: { title: string; color: string } | null;
|
||||
} = {},
|
||||
) {
|
||||
@@ -224,7 +230,15 @@ export class RelayBridge {
|
||||
for (const tab of this.#tabs.values()) {
|
||||
tab.attached = false;
|
||||
tab.attaching = null;
|
||||
// The extension dissolves omp groups on disconnect (or died along
|
||||
// with them); grouping state is unknowable until the next hello.
|
||||
// Without this reset, the next hello's groupId=-1 snapshots would
|
||||
// read as the user dragging every tab out (permanent opt-out).
|
||||
tab.grouped = false;
|
||||
tab.grouping = false;
|
||||
tab.ompGroupId = undefined;
|
||||
}
|
||||
this.#groupQueue.length = 0;
|
||||
}
|
||||
|
||||
extMessage(socket: RelaySocket, raw: string): void {
|
||||
@@ -314,6 +328,14 @@ export class RelayBridge {
|
||||
const touched = new Set<number>();
|
||||
for (const ref of conn.sessions.values()) touched.add(ref.tabId);
|
||||
conn.sessions.clear();
|
||||
// Tabs this client claimed leave the omp group unless another claimant
|
||||
// remains — session holders don't count: the long-lived registry
|
||||
// connection holds sessions on every tab without driving any of them.
|
||||
for (const tabId of conn.claims) {
|
||||
const tab = this.#tabs.get(tabId);
|
||||
if (tab) this.#syncTabGrouping(tab);
|
||||
}
|
||||
conn.claims.clear();
|
||||
// Drop the debugger (and its infobar) from tabs nobody drives anymore.
|
||||
for (const tabId of touched) {
|
||||
if (this.#sessionHolders(tabId).length > 0) continue;
|
||||
@@ -377,6 +399,13 @@ export class RelayBridge {
|
||||
this.#reply(conn, msg, {});
|
||||
return;
|
||||
}
|
||||
// Relay-private claim: the omp tab worker marks the page it was spawned
|
||||
// to drive. Never forwarded — real Chrome rejects the unknown method.
|
||||
if (msg.method === "OMP.claimTarget") {
|
||||
this.#claimTab(conn, tabId);
|
||||
this.#reply(conn, msg, {});
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const result = await this.#rpc({
|
||||
op: "send",
|
||||
@@ -391,6 +420,30 @@ export class RelayBridge {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Record `conn` as a driver of the tab and reconcile grouping. Claims are
|
||||
* explicit (worker adoption or tab creation) rather than inferred from
|
||||
* command traffic: target discovery scans every page with the same
|
||||
* commands a driver sends, so inference would sweep all tabs.
|
||||
*/
|
||||
#claimTab(conn: CdpConnection, tabId: number): void {
|
||||
const tab = this.#tabs.get(tabId);
|
||||
if (!tab) return;
|
||||
if (!conn.claims.has(tabId)) {
|
||||
conn.claims.add(tabId);
|
||||
this.#log("tab claimed", { conn: conn.id, tabId });
|
||||
}
|
||||
this.#syncTabGrouping(tab);
|
||||
}
|
||||
|
||||
/** True while any downstream connection claims the tab as its drive target. */
|
||||
#claimed(tabId: number): boolean {
|
||||
for (const conn of this.#conns.values()) {
|
||||
if (conn.claims.has(tabId)) return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/** Tab pseudo-sessions only exist to satisfy puppeteer's Target hierarchy. */
|
||||
#handleTabSessionCommand(conn: CdpConnection, msg: CdpCommand, ref: SessionRef): void {
|
||||
switch (msg.method) {
|
||||
@@ -500,6 +553,8 @@ export class RelayBridge {
|
||||
typeof msg.params?.url === "string" && msg.params.url.length > 0 ? msg.params.url : "about:blank";
|
||||
const result = (await this.#rpc({ op: "createTab", url })) as { tab: TabSnapshot };
|
||||
this.#onTabUpsert(result.tab);
|
||||
// Creating a tab is an explicit act of driving it.
|
||||
this.#claimTab(conn, result.tab.tabId);
|
||||
this.#reply(conn, msg, { targetId: pageTargetId(result.tab.tabId) });
|
||||
return;
|
||||
}
|
||||
@@ -605,6 +660,9 @@ export class RelayBridge {
|
||||
tab.attached = false;
|
||||
tab.attaching = null;
|
||||
tab.banned = true;
|
||||
// The user dismissed the debugger infobar (or the attach was torn
|
||||
// down): release the tab's omp-group membership too.
|
||||
this.#syncTabGrouping(tab);
|
||||
this.#retractTab(tab);
|
||||
}
|
||||
|
||||
@@ -613,6 +671,7 @@ export class RelayBridge {
|
||||
if (!tab) return;
|
||||
this.#retractTab(tab);
|
||||
this.#tabs.delete(tabId);
|
||||
for (const conn of this.#conns.values()) conn.claims.delete(tabId);
|
||||
}
|
||||
|
||||
#onTabUpsert(snap: TabSnapshot, opts: { silent?: boolean } = {}): void {
|
||||
@@ -632,7 +691,7 @@ export class RelayBridge {
|
||||
}
|
||||
if (opts.silent) return;
|
||||
const eligible = this.#eligible(tab);
|
||||
this.#syncTabGrouping(tab, eligible);
|
||||
this.#syncTabGrouping(tab);
|
||||
if (eligible && !tab.announced) {
|
||||
tab.announced = true;
|
||||
for (const conn of this.#conns.values()) {
|
||||
@@ -663,13 +722,13 @@ export class RelayBridge {
|
||||
|
||||
// ---- tab grouping -----------------------------------------------------------
|
||||
|
||||
/** A tab belongs in the omp group when controllable, unpinned, not user-opted-out, and not already in a user group. */
|
||||
/** A tab belongs in the omp group when claimed by a client, controllable, unpinned, not user-opted-out, and not already in a user group. */
|
||||
#groupWorthy(tab: TabState): boolean {
|
||||
if (!this.#eligible(tab) || tab.pinned || tab.groupOptOut) return false;
|
||||
if (!this.#claimed(tab.tabId) || !this.#eligible(tab) || tab.pinned || tab.groupOptOut) return false;
|
||||
return tab.grouped || tab.groupId === -1;
|
||||
}
|
||||
|
||||
/** Group every currently worthy tab (extension hello / reconnect). */
|
||||
/** Re-group every claimed tab (extension hello / reconnect). */
|
||||
#syncGrouping(): void {
|
||||
if (!this.#group) return;
|
||||
const worthy = [...this.#tabs.values()].filter(tab => this.#groupWorthy(tab) && !tab.grouped && !tab.grouping);
|
||||
@@ -677,49 +736,68 @@ export class RelayBridge {
|
||||
}
|
||||
|
||||
/** Reconcile one tab's group membership after a lifecycle event. */
|
||||
#syncTabGrouping(tab: TabState, eligible: boolean): void {
|
||||
#syncTabGrouping(tab: TabState): void {
|
||||
if (!this.#group) return;
|
||||
if (eligible && this.#groupWorthy(tab) && !tab.grouped && !tab.grouping) {
|
||||
this.#requestGroup([tab]);
|
||||
if (this.#groupWorthy(tab)) {
|
||||
if (!tab.grouped && !tab.grouping) this.#requestGroup([tab]);
|
||||
return;
|
||||
}
|
||||
if (!eligible && tab.grouped) {
|
||||
if (tab.grouped) {
|
||||
tab.grouped = false;
|
||||
tab.ompGroupId = undefined;
|
||||
void this.#rpc({ op: "ungroup", tabIds: [tab.tabId] }).catch(() => {});
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Queue tabs for grouping and drain serially. Overlapping group RPCs race
|
||||
* the extension's non-atomic query→create→set-title sequence and mint
|
||||
* duplicate omp groups, so at most one group RPC is ever in flight.
|
||||
*/
|
||||
#requestGroup(tabs: TabState[]): void {
|
||||
if (!this.#group) return;
|
||||
for (const tab of tabs) {
|
||||
tab.grouping = true;
|
||||
this.#groupQueue.push(tab);
|
||||
}
|
||||
if (!this.#groupDraining) void this.#drainGroupQueue();
|
||||
}
|
||||
|
||||
async #drainGroupQueue(): Promise<void> {
|
||||
const group = this.#group;
|
||||
if (!group) return;
|
||||
const tabIds = tabs.map(tab => tab.tabId);
|
||||
for (const tab of tabs) tab.grouping = true;
|
||||
void this.#rpc({ op: "group", tabIds, title: group.title, color: group.color })
|
||||
.then(result => {
|
||||
// Extension replies { grouped: { [tabId]: groupId } }; validate per entry.
|
||||
const grouped: Record<string, unknown> =
|
||||
result &&
|
||||
typeof result === "object" &&
|
||||
"grouped" in result &&
|
||||
result.grouped &&
|
||||
typeof result.grouped === "object"
|
||||
? (result.grouped as Record<string, unknown>)
|
||||
: {};
|
||||
for (const tab of tabs) {
|
||||
const groupId = grouped[String(tab.tabId)];
|
||||
if (typeof groupId !== "number") continue;
|
||||
tab.grouped = true;
|
||||
tab.ompGroupId = groupId;
|
||||
this.#groupDraining = true;
|
||||
try {
|
||||
while (this.#groupQueue.length > 0) {
|
||||
const batch = this.#groupQueue.splice(0);
|
||||
const tabIds = batch.map(tab => tab.tabId);
|
||||
try {
|
||||
const result = await this.#rpc({ op: "group", tabIds, title: group.title, color: group.color });
|
||||
// Extension replies { grouped: { [tabId]: groupId } }; validate per entry.
|
||||
const grouped: Record<string, unknown> =
|
||||
result &&
|
||||
typeof result === "object" &&
|
||||
"grouped" in result &&
|
||||
result.grouped &&
|
||||
typeof result.grouped === "object"
|
||||
? (result.grouped as Record<string, unknown>)
|
||||
: {};
|
||||
for (const tab of batch) {
|
||||
const groupId = grouped[String(tab.tabId)];
|
||||
if (typeof groupId !== "number") continue;
|
||||
tab.grouped = true;
|
||||
tab.ompGroupId = groupId;
|
||||
}
|
||||
this.#log("grouped tabs", { tabIds, grouped });
|
||||
} catch (err) {
|
||||
this.#log("tab grouping failed", { error: err instanceof Error ? err.message : String(err) });
|
||||
} finally {
|
||||
for (const tab of batch) tab.grouping = false;
|
||||
}
|
||||
this.#log("grouped tabs", { tabIds, grouped });
|
||||
})
|
||||
.catch(err => {
|
||||
this.#log("tab grouping failed", { error: err instanceof Error ? err.message : String(err) });
|
||||
})
|
||||
.finally(() => {
|
||||
for (const tab of tabs) tab.grouping = false;
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
this.#groupDraining = false;
|
||||
}
|
||||
}
|
||||
|
||||
/** Tear a tab out of every downstream connection (closed, detached, or now ineligible). */
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
// packages/browser-relay/extension/background.ts
|
||||
// extension/background.ts
|
||||
var DEFAULT_PORT = 9224;
|
||||
var PING_INTERVAL_MS = 20000;
|
||||
var RECONNECT_MIN_MS = 1000;
|
||||
@@ -28,8 +28,15 @@ function snapshot(tab) {
|
||||
};
|
||||
}
|
||||
var ompGroupTitle = null;
|
||||
var groupOps = Promise.resolve();
|
||||
function enqueueGroupOp(fn) {
|
||||
const result = groupOps.then(fn, fn);
|
||||
groupOps = result.catch(() => {});
|
||||
return result;
|
||||
}
|
||||
async function groupTabs(tabIds, title, color) {
|
||||
ompGroupTitle = title;
|
||||
chrome.storage.session.set({ ompGroupTitle: title });
|
||||
const byWindow = new Map;
|
||||
for (const tabId of tabIds) {
|
||||
try {
|
||||
@@ -44,7 +51,19 @@ async function groupTabs(tabIds, title, color) {
|
||||
const grouped = {};
|
||||
for (const [windowId, ids] of byWindow) {
|
||||
const existing = await chrome.tabGroups.query({ title, windowId });
|
||||
const groupId = await chrome.tabs.group(existing[0] ? { tabIds: ids, groupId: existing[0].id } : { tabIds: ids });
|
||||
let groupId;
|
||||
if (existing[0]) {
|
||||
groupId = existing[0].id;
|
||||
for (const dupe of existing.slice(1)) {
|
||||
const dupeTabs = await chrome.tabs.query({ groupId: dupe.id });
|
||||
const dupeIds = dupeTabs.map((tab) => tab.id).filter((id) => id !== undefined);
|
||||
if (dupeIds.length > 0)
|
||||
await chrome.tabs.group({ tabIds: dupeIds, groupId });
|
||||
}
|
||||
await chrome.tabs.group({ tabIds: ids, groupId });
|
||||
} else {
|
||||
groupId = await chrome.tabs.group({ tabIds: ids });
|
||||
}
|
||||
await chrome.tabGroups.update(groupId, { title, color });
|
||||
for (const id of ids)
|
||||
grouped[String(id)] = groupId;
|
||||
@@ -52,6 +71,10 @@ async function groupTabs(tabIds, title, color) {
|
||||
return { grouped };
|
||||
}
|
||||
async function restoreGroups() {
|
||||
if (!ompGroupTitle) {
|
||||
const stored = await chrome.storage.session.get({ ompGroupTitle: "" }).catch(() => ({ ompGroupTitle: "" }));
|
||||
ompGroupTitle = typeof stored.ompGroupTitle === "string" && stored.ompGroupTitle ? stored.ompGroupTitle : null;
|
||||
}
|
||||
if (!ompGroupTitle)
|
||||
return;
|
||||
const groups = await chrome.tabGroups.query({ title: ompGroupTitle }).catch(() => []);
|
||||
@@ -121,9 +144,9 @@ async function runRpc(msg) {
|
||||
return {};
|
||||
}
|
||||
case "group":
|
||||
return await groupTabs(msg.tabIds, msg.title, msg.color);
|
||||
return await enqueueGroupOp(() => groupTabs(msg.tabIds, msg.title, msg.color));
|
||||
case "ungroup":
|
||||
await chrome.tabs.ungroup(msg.tabIds).catch(() => {});
|
||||
await enqueueGroupOp(() => chrome.tabs.ungroup(msg.tabIds).catch(() => {}));
|
||||
return {};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -19,7 +19,7 @@ export interface RelayServerOptions {
|
||||
port: number;
|
||||
/** Shared secret the extension must present as `?token=`; unset disables the check. */
|
||||
token?: string;
|
||||
/** Group controllable tabs under one per-window Chrome tab group (default on); `false` disables. */
|
||||
/** Group tabs the agent actively drives under one per-window Chrome tab group (default on); `false` disables. */
|
||||
group?: boolean | { title: string; color: string };
|
||||
log?: (message: string, data?: Record<string, unknown>) => void;
|
||||
}
|
||||
|
||||
@@ -20,6 +20,14 @@ import { JsRuntime, type RuntimeHooks } from "../../eval/js/shared/runtime";
|
||||
import { resizeImage } from "../../utils/image-resize";
|
||||
import { resolveToCwd } from "../path-utils";
|
||||
import { formatScreenshot } from "../render-utils";
|
||||
import {
|
||||
bindRunFacade,
|
||||
CELL_BUDGET_SLACK_MS,
|
||||
markHandled,
|
||||
resolvePredicateTimeout,
|
||||
type WaitPredicateOptions,
|
||||
waitForRun,
|
||||
} from "../run-scope";
|
||||
import { ToolAbortError, ToolError, throwIfAborted } from "../tool-errors";
|
||||
import {
|
||||
type AriaSnapshotOptions,
|
||||
@@ -36,14 +44,6 @@ import {
|
||||
loadPuppeteerInWorker,
|
||||
} from "./launch";
|
||||
import { extractReadableFromHtml, type ReadableFormat } from "./readable";
|
||||
import {
|
||||
bindBrowserRunFacade,
|
||||
CELL_BUDGET_SLACK_MS,
|
||||
markHandled,
|
||||
resolvePredicateTimeout,
|
||||
type WaitPredicateOptions,
|
||||
waitForBrowserRun,
|
||||
} from "./run-cancellation";
|
||||
import { cloneSafe, RunOutput } from "./run-output";
|
||||
import type {
|
||||
Observation,
|
||||
@@ -827,6 +827,7 @@ export class WorkerCore {
|
||||
const page = await target.page();
|
||||
if (!page) throw new ToolError(`Target ${payload.targetId} is no longer available on the attached browser`);
|
||||
this.#page = page;
|
||||
await this.#claimRelayTarget(page);
|
||||
this.#observeDialogs();
|
||||
if (payload.dialogs) this.#applyDialogPolicy(payload.dialogs);
|
||||
if (payload.url) {
|
||||
@@ -853,6 +854,26 @@ export class WorkerCore {
|
||||
throw new ToolError(`Target ${targetId} is no longer available on the attached browser`);
|
||||
}
|
||||
|
||||
/**
|
||||
* Tell the omp browser relay this worker drives the adopted page, so the
|
||||
* relay adds it to the per-window "omp" tab group. Best-effort: plain CDP
|
||||
* backends (real Chrome, cmux) reject the relay-private method.
|
||||
*/
|
||||
async #claimRelayTarget(page: Page): Promise<void> {
|
||||
let session: CDPSession | undefined;
|
||||
try {
|
||||
session = await page.createCDPSession();
|
||||
// Puppeteer's protocol map cannot express the relay-private method; the
|
||||
// send signature is otherwise identical.
|
||||
const raw = session as unknown as { send(method: string): Promise<unknown> };
|
||||
await raw.send("OMP.claimTarget");
|
||||
} catch {
|
||||
// Not the omp relay; nothing to claim.
|
||||
} finally {
|
||||
await session?.detach().catch(() => undefined);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Best-effort unblocking of a wedged target during post-timeout recovery: dismiss any
|
||||
* open JS dialog and stop a pending navigation over a raw CDP session (created on the
|
||||
@@ -975,9 +996,9 @@ export class WorkerCore {
|
||||
const runtime = this.#ensureRuntime(msg.session);
|
||||
runtime.setCwd(msg.session.cwd);
|
||||
runtime.setRunScope({
|
||||
page: bindBrowserRunFacade(runPage.page, signal),
|
||||
browser: bindBrowserRunFacade(browser, signal),
|
||||
tab: bindBrowserRunFacade(tabApi, signal),
|
||||
page: bindRunFacade(runPage.page, signal),
|
||||
browser: bindRunFacade(browser, signal),
|
||||
tab: bindRunFacade(tabApi, signal),
|
||||
assert: (cond: unknown, text?: string): void => {
|
||||
if (!cond) throw new ToolError(text ?? "Assertion failed");
|
||||
},
|
||||
@@ -991,7 +1012,7 @@ export class WorkerCore {
|
||||
: { timeout: resolvePredicateTimeout(msg.timeoutMs, opts?.timeout), interval: opts?.interval };
|
||||
return markHandled(
|
||||
this.#runOp(active, label, signal, Number.POSITIVE_INFINITY, sig =>
|
||||
waitForBrowserRun(msOrPredicate, sig, resolved),
|
||||
waitForRun(msOrPredicate, sig, resolved),
|
||||
),
|
||||
);
|
||||
},
|
||||
|
||||
+10
-10
@@ -1,10 +1,10 @@
|
||||
import { untilAborted } from "@oh-my-pi/pi-utils";
|
||||
import { ToolError, throwIfAborted } from "../tool-errors";
|
||||
import { ToolError, throwIfAborted } from "./tool-errors";
|
||||
|
||||
/**
|
||||
* Marks a run-scoped promise as observed without changing its behavior for awaited callers.
|
||||
*
|
||||
* Browser run teardown aborts can reject promises created for evaluated code after user code
|
||||
* Run teardown aborts can reject promises created for evaluated code after user code
|
||||
* has stopped observing them (for example fire-and-forget `wait()`/facade calls). In 16.3.0
|
||||
* those zero-consumer rejections reached the process-level `unhandledRejection` handler and
|
||||
* killed every subagent sharing the process (issues #4499/#4672). Attaching a no-op rejection
|
||||
@@ -33,8 +33,8 @@ export interface WaitPredicateOptions {
|
||||
/**
|
||||
* Effective `wait(predicate)` deadline for a given cell budget. Always strictly below
|
||||
* the cell budget so the named `wait(predicate) timed out` error wins the race against
|
||||
* the opaque whole-cell "Browser code execution timed out". `0`/`Infinity` ("disable")
|
||||
* map to the largest bounded deadline; negative/NaN garbage falls back to the default.
|
||||
* the opaque whole-cell execution timeout. `0`/`Infinity` ("disable") map to the largest
|
||||
* bounded deadline; negative/NaN garbage falls back to the default.
|
||||
*/
|
||||
export function resolvePredicateTimeout(cellTimeoutMs: number, explicit?: number): number {
|
||||
const budgetBound = Math.max(1, cellTimeoutMs - CELL_BUDGET_SLACK_MS);
|
||||
@@ -44,15 +44,15 @@ export function resolvePredicateTimeout(cellTimeoutMs: number, explicit?: number
|
||||
}
|
||||
|
||||
/**
|
||||
* Run-scoped `wait()` helper for evaluated browser code, honoring the owning run's
|
||||
* cancellation signal.
|
||||
* Run-scoped `wait()` helper for evaluated code (browser and computer workers), honoring
|
||||
* the owning run's cancellation signal.
|
||||
*
|
||||
* - `wait(ms)` sleeps for `ms` milliseconds.
|
||||
* - `wait(fn, { timeout?, interval? })` polls `fn` (sync or async) until it returns a
|
||||
* truthy value and resolves with that value; throws a named `ToolError` on timeout
|
||||
* instead of stalling into the whole-cell deadline. Predicate errors propagate.
|
||||
*/
|
||||
export function waitForBrowserRun(
|
||||
export function waitForRun(
|
||||
msOrPredicate: number | (() => unknown),
|
||||
signal: AbortSignal,
|
||||
opts?: WaitPredicateOptions,
|
||||
@@ -86,8 +86,8 @@ export function waitForBrowserRun(
|
||||
return markHandled(promise);
|
||||
}
|
||||
|
||||
/** Binds a long-lived browser facade to one evaluated run's abort signal. */
|
||||
export function bindBrowserRunFacade<T extends object>(target: T, signal: AbortSignal): T {
|
||||
/** Binds a long-lived scope facade (page/tab/desktop objects) to one evaluated run's abort signal. */
|
||||
export function bindRunFacade<T extends object>(target: T, signal: AbortSignal): T {
|
||||
const cache = new Map<PropertyKey, unknown>();
|
||||
return new Proxy(target, {
|
||||
get(current, prop) {
|
||||
@@ -121,7 +121,7 @@ export function bindBrowserRunFacade<T extends object>(target: T, signal: AbortS
|
||||
// brand-check internal slots that a Proxy cannot forward, and reading a
|
||||
// signal needs no abort gating anyway.
|
||||
if (value instanceof AbortSignal) return value;
|
||||
const wrapped = bindBrowserRunFacade(value, signal);
|
||||
const wrapped = bindRunFacade(value, signal);
|
||||
cache.set(prop, wrapped);
|
||||
return wrapped;
|
||||
}
|
||||
@@ -23,7 +23,7 @@
|
||||
* immediately instead of blocking to the run's timeout.
|
||||
* 2. When the in-flight run is doing work that does NOT make another cmux
|
||||
* socket request (e.g. `await wait(60_000)`), releasing the tab still
|
||||
* unwinds the run — proving `closeAc.signal` reaches `waitForBrowserRun`
|
||||
* unwinds the run — proving `closeAc.signal` reaches `waitForRun`
|
||||
* and the facade proxies, not just the outer race. (Reviewer feedback
|
||||
* from PR #4502.)
|
||||
*/
|
||||
@@ -218,7 +218,7 @@ describe("browser tab-supervisor — cmux tab close mid-run (#4499)", () => {
|
||||
|
||||
const session = makeSession("/tmp");
|
||||
// The user code awaits `wait(60_000)` — which drives
|
||||
// `waitForBrowserRun(60_000, signal)` -> `untilAborted(signal,
|
||||
// `waitForRun(60_000, signal)` -> `untilAborted(signal,
|
||||
// () => Bun.sleep(60_000))` INSIDE the runtime. Nothing hits the
|
||||
// cmux socket, so on `main` the reviewer's exact scenario applies:
|
||||
// even after `pending.reject` unblocks the caller, `runCmuxCode`
|
||||
@@ -251,7 +251,7 @@ describe("browser tab-supervisor — cmux tab close mid-run (#4499)", () => {
|
||||
// the map. This is the wire the reviewer asked us to check: the
|
||||
// tab-close event must reach the cmux run body, not only the
|
||||
// awaiting caller. Its `.signal.aborted` is the observable proof
|
||||
// that `waitForBrowserRun` / cmux socket calls will unwind
|
||||
// that `waitForRun` / cmux socket calls will unwind
|
||||
// synchronously (via `untilAborted`) instead of blocking to the
|
||||
// 60_000ms timeout.
|
||||
const pendingBeforeRelease = [...(tabBeforeRelease?.pending.values() ?? [])];
|
||||
|
||||
@@ -1,17 +1,44 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { RelayBridge, type RelaySocket } from "@oh-my-pi/pi-coding-agent/tools/browser/relay/bridge";
|
||||
import type { RelayToExtMessage, TabSnapshot } from "@oh-my-pi/pi-coding-agent/tools/browser/relay/protocol";
|
||||
import type {
|
||||
RelayRpcRequest,
|
||||
RelayToExtMessage,
|
||||
TabSnapshot,
|
||||
} from "@oh-my-pi/pi-coding-agent/tools/browser/relay/protocol";
|
||||
|
||||
/** A relay→extension RPC narrowed to one op, tabIds/title/etc. included. */
|
||||
type ExtRpc<Op extends RelayRpcRequest["op"]> = { t: "rpc"; id: number } & Extract<RelayRpcRequest, { op: Op }>;
|
||||
|
||||
class FakeExtSocket implements RelaySocket {
|
||||
readonly messages: RelayToExtMessage[] = [];
|
||||
readonly #acked = new Set<number>();
|
||||
send(text: string): void {
|
||||
this.messages.push(JSON.parse(text) as RelayToExtMessage);
|
||||
}
|
||||
close(): void {}
|
||||
rpcs(op: string): Array<Extract<RelayToExtMessage, { t: "rpc" }>> {
|
||||
return this.messages.filter(
|
||||
(msg): msg is Extract<RelayToExtMessage, { t: "rpc" }> => msg.t === "rpc" && msg.op === op,
|
||||
);
|
||||
rpcs<Op extends RelayRpcRequest["op"]>(op: Op): Array<ExtRpc<Op>> {
|
||||
return this.messages.filter((msg): msg is ExtRpc<Op> => msg.t === "rpc" && msg.op === op);
|
||||
}
|
||||
/** RPC requests of `op` not yet answered through {@link ack}. */
|
||||
pending<Op extends RelayRpcRequest["op"]>(op: Op): Array<ExtRpc<Op>> {
|
||||
return this.rpcs(op).filter(msg => !this.#acked.has(msg.id));
|
||||
}
|
||||
markAcked(id: number): void {
|
||||
this.#acked.add(id);
|
||||
}
|
||||
}
|
||||
|
||||
/** Downstream puppeteer-side socket capturing bridge emissions. */
|
||||
class FakeCdpSocket implements RelaySocket {
|
||||
readonly messages: Array<Record<string, unknown>> = [];
|
||||
send(text: string): void {
|
||||
this.messages.push(JSON.parse(text) as Record<string, unknown>);
|
||||
}
|
||||
close(): void {}
|
||||
sessionFor(commandId: number): string | undefined {
|
||||
const msg = this.messages.find(m => m.id === commandId);
|
||||
const result = msg && "result" in msg && msg.result && typeof msg.result === "object" ? msg.result : undefined;
|
||||
return result && "sessionId" in result && typeof result.sessionId === "string" ? result.sessionId : undefined;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -41,63 +68,215 @@ function connect(bridge: RelayBridge, socket: FakeExtSocket, tabs: TabSnapshot[]
|
||||
);
|
||||
}
|
||||
|
||||
/** Answer every unanswered extension RPC of `op` with `ok: true` and `result`. */
|
||||
function ack(bridge: RelayBridge, socket: FakeExtSocket, op: RelayRpcRequest["op"], result: unknown = {}): void {
|
||||
for (const rpc of socket.pending(op)) {
|
||||
socket.markAcked(rpc.id);
|
||||
bridge.extMessage(socket, JSON.stringify({ t: "rpcResult", id: rpc.id, ok: true, result }));
|
||||
}
|
||||
}
|
||||
|
||||
/** Flush the rpc .then() microtask chains (no timers involved). */
|
||||
async function flush(): Promise<void> {
|
||||
for (let i = 0; i < 5; i++) await Promise.resolve();
|
||||
}
|
||||
|
||||
let msgSeq = 100;
|
||||
|
||||
/** Attach to a tab's page target and return the minted page session id. */
|
||||
async function attachPage(
|
||||
bridge: RelayBridge,
|
||||
ext: FakeExtSocket,
|
||||
cdp: FakeCdpSocket,
|
||||
connId: number,
|
||||
tabId: number,
|
||||
): Promise<string> {
|
||||
const attachId = ++msgSeq;
|
||||
bridge.cdpMessage(
|
||||
connId,
|
||||
JSON.stringify({
|
||||
id: attachId,
|
||||
method: "Target.attachToTarget",
|
||||
params: { targetId: `PAGE${tabId}`, flatten: true },
|
||||
}),
|
||||
);
|
||||
ack(bridge, ext, "attach");
|
||||
await flush();
|
||||
const sessionId = cdp.sessionFor(attachId);
|
||||
if (!sessionId) throw new Error(`attachToTarget for tab ${tabId} did not produce a session`);
|
||||
return sessionId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Emulate the omp tab worker adopting a tab: attach to its page target, then
|
||||
* claim it as this connection's drive target.
|
||||
*/
|
||||
async function claimTab(
|
||||
bridge: RelayBridge,
|
||||
ext: FakeExtSocket,
|
||||
cdp: FakeCdpSocket,
|
||||
connId: number,
|
||||
tabId: number,
|
||||
): Promise<void> {
|
||||
const sessionId = await attachPage(bridge, ext, cdp, connId, tabId);
|
||||
bridge.cdpMessage(connId, JSON.stringify({ id: ++msgSeq, sessionId, method: "OMP.claimTarget" }));
|
||||
await flush();
|
||||
}
|
||||
|
||||
describe("RelayBridge tab grouping", () => {
|
||||
it("groups only controllable, unpinned, ungrouped tabs on hello", () => {
|
||||
it("groups nothing on hello or tab lifecycle events — only claimed tabs join the omp group", () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const socket = new FakeExtSocket();
|
||||
connect(bridge, socket, [
|
||||
tab({ tabId: 1 }),
|
||||
tab({ tabId: 2, url: "chrome://settings/" }),
|
||||
tab({ tabId: 3, pinned: true }),
|
||||
tab({ tabId: 4, groupId: 77 }), // already in a user group
|
||||
tab({ tabId: 5, url: "about:blank" }),
|
||||
]);
|
||||
const groups = socket.rpcs("group");
|
||||
expect(groups).toHaveLength(1);
|
||||
const group = groups[0]! as { tabIds: number[]; title: string; color: string };
|
||||
expect(group.tabIds.toSorted()).toEqual([1, 5]);
|
||||
expect(group.title).toBe("omp");
|
||||
expect(group.color).toBe("cyan");
|
||||
connect(bridge, socket, [tab({ tabId: 1 }), tab({ tabId: 2 }), tab({ tabId: 3, url: "about:blank" })]);
|
||||
bridge.extMessage(socket, JSON.stringify({ t: "tabCreated", tab: tab({ tabId: 9 }) }));
|
||||
expect(socket.rpcs("group")).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("does not issue group RPCs when grouping is disabled", () => {
|
||||
it("never groups from command traffic: a discovery scan sending page commands to every tab is not driving", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 1 }), tab({ tabId: 2 })]);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
// pickElectronTarget materializes every discovered page, which makes
|
||||
// puppeteer send Page.enable/Page.getFrameTree to all of them.
|
||||
for (const tabId of [1, 2]) {
|
||||
const sessionId = await attachPage(bridge, ext, cdp, connId, tabId);
|
||||
bridge.cdpMessage(connId, JSON.stringify({ id: ++msgSeq, sessionId, method: "Page.enable" }));
|
||||
bridge.cdpMessage(connId, JSON.stringify({ id: ++msgSeq, sessionId, method: "Page.getFrameTree" }));
|
||||
}
|
||||
await flush();
|
||||
expect(ext.rpcs("group")).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("groups exactly the tab a client claims", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 1 }), tab({ tabId: 2 })]);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
await claimTab(bridge, ext, cdp, connId, 1);
|
||||
const groups = ext.rpcs("group");
|
||||
expect(groups).toHaveLength(1);
|
||||
expect(groups[0]!.tabIds).toEqual([1]);
|
||||
expect(groups[0]!.title).toBe("omp");
|
||||
expect(groups[0]!.color).toBe("cyan");
|
||||
});
|
||||
|
||||
it("never groups pinned tabs or tabs in a user group, even when claimed", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 3, pinned: true }), tab({ tabId: 4, groupId: 77 })]);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
await claimTab(bridge, ext, cdp, connId, 3);
|
||||
await claimTab(bridge, ext, cdp, connId, 4);
|
||||
expect(ext.rpcs("group")).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("does not issue group RPCs when grouping is disabled", async () => {
|
||||
const bridge = new RelayBridge({});
|
||||
const socket = new FakeExtSocket();
|
||||
connect(bridge, socket, [tab({ tabId: 1 })]);
|
||||
expect(socket.rpcs("group")).toHaveLength(0);
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 1 })]);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
await claimTab(bridge, ext, cdp, connId, 1);
|
||||
expect(ext.rpcs("group")).toHaveLength(0);
|
||||
});
|
||||
|
||||
it("auto-claims a tab created through Target.createTarget", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, []);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
bridge.cdpMessage(
|
||||
connId,
|
||||
JSON.stringify({ id: ++msgSeq, method: "Target.createTarget", params: { url: "https://example.com/" } }),
|
||||
);
|
||||
ack(bridge, ext, "createTab", { tab: tab({ tabId: 9 }) });
|
||||
await flush();
|
||||
const groups = ext.rpcs("group");
|
||||
expect(groups).toHaveLength(1);
|
||||
expect(groups[0]!.tabIds).toEqual([9]);
|
||||
});
|
||||
|
||||
it("never re-groups a tab the user pulled out of the omp group", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const socket = new FakeExtSocket();
|
||||
connect(bridge, socket, [tab({ tabId: 1 })]);
|
||||
const first = socket.rpcs("group")[0]!;
|
||||
bridge.extMessage(
|
||||
socket,
|
||||
JSON.stringify({ t: "rpcResult", id: first.id, ok: true, result: { grouped: { "1": 42 } } }),
|
||||
);
|
||||
// Flush the rpc .then() microtask chain (no timers involved).
|
||||
await Promise.resolve();
|
||||
await Promise.resolve();
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 1 })]);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
await claimTab(bridge, ext, cdp, connId, 1);
|
||||
ack(bridge, ext, "group", { grouped: { "1": 42 } });
|
||||
await flush();
|
||||
// Chrome reports the grouping we just made — no opt-out.
|
||||
bridge.extMessage(socket, JSON.stringify({ t: "tabUpdated", tab: tab({ tabId: 1, groupId: 42 }) }));
|
||||
bridge.extMessage(ext, JSON.stringify({ t: "tabUpdated", tab: tab({ tabId: 1, groupId: 42 }) }));
|
||||
// The user drags the tab out of the group.
|
||||
bridge.extMessage(socket, JSON.stringify({ t: "tabUpdated", tab: tab({ tabId: 1, groupId: -1 }) }));
|
||||
// A later navigation on the now-ungrouped tab must not re-group it.
|
||||
bridge.extMessage(ext, JSON.stringify({ t: "tabUpdated", tab: tab({ tabId: 1, groupId: -1 }) }));
|
||||
// A later navigation on the still-claimed tab must not re-group it.
|
||||
bridge.extMessage(
|
||||
socket,
|
||||
ext,
|
||||
JSON.stringify({ t: "tabUpdated", tab: tab({ tabId: 1, groupId: -1, url: "https://example.com/other" }) }),
|
||||
);
|
||||
expect(socket.rpcs("group")).toHaveLength(1);
|
||||
expect(ext.rpcs("group")).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("groups a newly created controllable tab", () => {
|
||||
it("ungroups when the claiming client disconnects, even while another connection still holds sessions", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const socket = new FakeExtSocket();
|
||||
connect(bridge, socket, []);
|
||||
bridge.extMessage(socket, JSON.stringify({ t: "tabCreated", tab: tab({ tabId: 9 }) }));
|
||||
const groups = socket.rpcs("group");
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 1 })]);
|
||||
// Long-lived registry connection: holds a session on the tab, never claims it.
|
||||
const registry = new FakeCdpSocket();
|
||||
const registryConn = bridge.cdpConnected(registry);
|
||||
await attachPage(bridge, ext, registry, registryConn, 1);
|
||||
// Worker connection: claims the tab.
|
||||
const worker = new FakeCdpSocket();
|
||||
const workerConn = bridge.cdpConnected(worker);
|
||||
await claimTab(bridge, ext, worker, workerConn, 1);
|
||||
ack(bridge, ext, "group", { grouped: { "1": 42 } });
|
||||
await flush();
|
||||
bridge.cdpClosed(workerConn);
|
||||
const ungroups = ext.rpcs("ungroup");
|
||||
expect(ungroups).toHaveLength(1);
|
||||
expect(ungroups[0]!.tabIds).toEqual([1]);
|
||||
});
|
||||
|
||||
it("never overlaps group RPCs: a tab claimed mid-flight waits for the pending group", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 1 }), tab({ tabId: 2 })]);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
await claimTab(bridge, ext, cdp, connId, 1);
|
||||
expect(ext.rpcs("group")).toHaveLength(1);
|
||||
// Concurrent group RPCs race Chrome's non-atomic query→create→set-title
|
||||
// and mint duplicate "omp" groups; the second request must queue.
|
||||
await claimTab(bridge, ext, cdp, connId, 2);
|
||||
expect(ext.rpcs("group")).toHaveLength(1);
|
||||
ack(bridge, ext, "group", { grouped: { "1": 42 } });
|
||||
await flush();
|
||||
const groups = ext.rpcs("group");
|
||||
expect(groups).toHaveLength(2);
|
||||
expect(groups[1]!.tabIds).toEqual([2]);
|
||||
});
|
||||
|
||||
it("regroups claimed tabs after an extension reconnect instead of treating the dissolve as user opt-out", async () => {
|
||||
const bridge = new RelayBridge({ group: { title: "omp", color: "cyan" } });
|
||||
const ext = new FakeExtSocket();
|
||||
connect(bridge, ext, [tab({ tabId: 1 })]);
|
||||
const cdp = new FakeCdpSocket();
|
||||
const connId = bridge.cdpConnected(cdp);
|
||||
await claimTab(bridge, ext, cdp, connId, 1);
|
||||
ack(bridge, ext, "group", { grouped: { "1": 42 } });
|
||||
await flush();
|
||||
// Relay/extension link drops: the extension dissolves the omp group on
|
||||
// disconnect, so the next hello reports groupId -1 for every tab.
|
||||
bridge.extClosed(ext);
|
||||
const ext2 = new FakeExtSocket();
|
||||
connect(bridge, ext2, [tab({ tabId: 1, groupId: -1 })]);
|
||||
const groups = ext2.rpcs("group");
|
||||
expect(groups).toHaveLength(1);
|
||||
expect((groups[0] as unknown as { tabIds: number[] }).tabIds).toEqual([9]);
|
||||
expect(groups[0]!.tabIds).toEqual([1]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
|
||||
import { postmortem } from "@oh-my-pi/pi-utils";
|
||||
import { JsRuntime, type RuntimeHooks } from "../../src/eval/js/shared/runtime";
|
||||
import { bindBrowserRunFacade, markHandled, waitForBrowserRun } from "../../src/tools/browser/run-cancellation";
|
||||
import { bindRunFacade, markHandled, waitForRun } from "../../src/tools/run-scope";
|
||||
import { ToolAbortError } from "../../src/tools/tool-errors";
|
||||
|
||||
async function collectUnhandledRejections(action: () => void | Promise<void>): Promise<unknown[]> {
|
||||
@@ -42,7 +42,7 @@ describe("browser run cancellation", () => {
|
||||
it("resolves run-scoped wait when the run is not aborted", async () => {
|
||||
const controller = new AbortController();
|
||||
|
||||
const wait = waitForBrowserRun(25, controller.signal);
|
||||
const wait = waitForRun(25, controller.signal);
|
||||
vi.advanceTimersByTime(25);
|
||||
|
||||
await expect(wait).resolves.toBeUndefined();
|
||||
@@ -50,7 +50,7 @@ describe("browser run cancellation", () => {
|
||||
|
||||
it("rejects run-scoped wait when the run aborts mid-sleep", async () => {
|
||||
const controller = new AbortController();
|
||||
const wait = waitForBrowserRun(1000, controller.signal);
|
||||
const wait = waitForRun(1000, controller.signal);
|
||||
|
||||
controller.abort(new Error("browser run ended"));
|
||||
|
||||
@@ -62,7 +62,7 @@ describe("browser run cancellation", () => {
|
||||
const controller = new AbortController();
|
||||
let calls = 0;
|
||||
|
||||
const wait = waitForBrowserRun(() => (++calls >= 3 ? "ready" : null), controller.signal, { interval: 10 });
|
||||
const wait = waitForRun(() => (++calls >= 3 ? "ready" : null), controller.signal, { interval: 10 });
|
||||
|
||||
await expect(wait).resolves.toBe("ready");
|
||||
expect(calls).toBe(3);
|
||||
@@ -72,7 +72,7 @@ describe("browser run cancellation", () => {
|
||||
vi.useRealTimers();
|
||||
const controller = new AbortController();
|
||||
|
||||
const wait = waitForBrowserRun(() => false, controller.signal, { timeout: 50, interval: 10 });
|
||||
const wait = waitForRun(() => false, controller.signal, { timeout: 50, interval: 10 });
|
||||
|
||||
await expect(wait).rejects.toThrow("wait(predicate) timed out after 50ms");
|
||||
});
|
||||
@@ -81,7 +81,7 @@ describe("browser run cancellation", () => {
|
||||
vi.useRealTimers();
|
||||
const controller = new AbortController();
|
||||
|
||||
const wait = waitForBrowserRun(() => false, controller.signal, { timeout: 5000 });
|
||||
const wait = waitForRun(() => false, controller.signal, { timeout: 5000 });
|
||||
controller.abort(new Error("browser run ended"));
|
||||
|
||||
await expect(wait).rejects.toThrow("browser run ended");
|
||||
@@ -90,7 +90,7 @@ describe("browser run cancellation", () => {
|
||||
it("rejects wait() input that is neither milliseconds nor a predicate", async () => {
|
||||
const controller = new AbortController();
|
||||
|
||||
await expect(waitForBrowserRun("soon" as never, controller.signal)).rejects.toThrow(
|
||||
await expect(waitForRun("soon" as never, controller.signal)).rejects.toThrow(
|
||||
"wait(...) expects milliseconds (number) or a predicate function to poll",
|
||||
);
|
||||
});
|
||||
@@ -99,7 +99,7 @@ describe("browser run cancellation", () => {
|
||||
const controller = new AbortController();
|
||||
|
||||
const reasons = await collectUnhandledRejections(async () => {
|
||||
void waitForBrowserRun(1000, controller.signal);
|
||||
void waitForRun(1000, controller.signal);
|
||||
controller.abort(postmortem.markExpectedCleanupError(new Error("browser run ended")));
|
||||
});
|
||||
|
||||
@@ -109,7 +109,7 @@ describe("browser run cancellation", () => {
|
||||
it("does not emit unhandledRejection when an unawaited facade method settles after abort", async () => {
|
||||
const controller = new AbortController();
|
||||
const deferred = Promise.withResolvers<string>();
|
||||
const facade = bindBrowserRunFacade(
|
||||
const facade = bindRunFacade(
|
||||
{
|
||||
readTitle(): Promise<string> {
|
||||
return deferred.promise;
|
||||
@@ -130,7 +130,7 @@ describe("browser run cancellation", () => {
|
||||
it("rejects awaited facade method calls that settle after abort", async () => {
|
||||
const controller = new AbortController();
|
||||
const deferred = Promise.withResolvers<string>();
|
||||
const facade = bindBrowserRunFacade(
|
||||
const facade = bindRunFacade(
|
||||
{
|
||||
readTitle(): Promise<string> {
|
||||
return deferred.promise;
|
||||
@@ -162,8 +162,8 @@ describe("browser run cancellation", () => {
|
||||
once: true,
|
||||
});
|
||||
runtime.setRunScope({
|
||||
wait: (ms: number): Promise<unknown> => waitForBrowserRun(ms, signal),
|
||||
tab: bindBrowserRunFacade(
|
||||
wait: (ms: number): Promise<unknown> => waitForRun(ms, signal),
|
||||
tab: bindRunFacade(
|
||||
{
|
||||
goto: async (url: string): Promise<void> => {
|
||||
state.lateNavigation = url;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { describe, expect, it, vi } from "bun:test";
|
||||
import { resolvePredicateTimeout } from "@oh-my-pi/pi-coding-agent/tools/browser/run-cancellation";
|
||||
import { resolvePredicateTimeout } from "@oh-my-pi/pi-coding-agent/tools/run-scope";
|
||||
import {
|
||||
dispatchScroll,
|
||||
normalizeSelector,
|
||||
|
||||
Reference in New Issue
Block a user