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
84 changes: 83 additions & 1 deletion apps/app/src/app/lib/openwork-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,15 @@ import { isOpenworkGatewayRuntime } from "./gateway-runtime";
import { isDesktopRuntime } from "./runtime-env";
import type { ExecResult, OpencodeConfigFile, WorkspaceInfo, WorkspaceList } from "./desktop";
import type { DenOrgMarketplace, DenOrgPluginResolved, DenResourceSnapshot } from "./den-types";
import type { CloudImportedMarketplace, CloudImportedPlugin } from "../cloud/import-state";
import type { CloudImportedMarketplace, CloudImportedPlugin, CloudImportedProvider } from "../cloud/import-state";

export type OpenworkServerCapabilities = {
skills: { read: boolean; write: boolean; source: "openwork" | "opencode" };
plugins: { read: boolean; write: boolean };
mcp: { read: boolean; write: boolean };
commands: { read: boolean; write: boolean };
config: { read: boolean; write: boolean };
providerSync?: boolean;
sandbox?: { enabled: boolean; backend: "none" | "docker" | "container" };
proxy?: { opencode: boolean };
toolProviders?: {
Expand All @@ -41,6 +42,69 @@ export type OpenworkServerCapabilities = {
};
};

export type OpenworkCloudProviderSyncRun = {
status: "applied" | "noop" | "failed" | "no_session";
message?: string;
};

export type OpenworkCloudProviderSyncStatus = {
hasSession: boolean;
lastRun: { at: string | number; status: OpenworkCloudProviderSyncRun["status"]; message?: string } | null;
providers: CloudImportedProvider[];
};

function parseCloudProviderSyncRun(value: unknown): OpenworkCloudProviderSyncRun {
if (!value || typeof value !== "object" || !("status" in value)) throw new Error("Invalid cloud provider sync response.");
const status = value.status;
if (status !== "applied" && status !== "noop" && status !== "failed" && status !== "no_session") {
throw new Error("Invalid cloud provider sync status.");
}
const message = "message" in value && typeof value.message === "string" ? value.message : undefined;
return { status, message };
}

function parseCloudImportedProvider(value: unknown): CloudImportedProvider | null {
if (!value || typeof value !== "object") return null;
if (
!("cloudProviderId" in value) || typeof value.cloudProviderId !== "string" ||
!("providerId" in value) || typeof value.providerId !== "string" ||
!("sourceProviderId" in value) || typeof value.sourceProviderId !== "string" ||
!("name" in value) || typeof value.name !== "string" ||
!("modelIds" in value) || !Array.isArray(value.modelIds) || !value.modelIds.every((item) => typeof item === "string")
) return null;
return {
cloudProviderId: value.cloudProviderId,
providerId: value.providerId,
sourceProviderId: value.sourceProviderId,
name: value.name,
source: "source" in value && typeof value.source === "string" ? value.source : null,
updatedAt: "updatedAt" in value && typeof value.updatedAt === "string" ? value.updatedAt : null,
modelIds: value.modelIds,
importedAt: "importedAt" in value && typeof value.importedAt === "number" ? value.importedAt : null,
};
}

function parseCloudProviderSyncStatus(value: unknown): OpenworkCloudProviderSyncStatus {
if (!value || typeof value !== "object" || !("hasSession" in value) || typeof value.hasSession !== "boolean" || !("providers" in value) || !Array.isArray(value.providers)) {
throw new Error("Invalid cloud provider sync status response.");
}
const providers: CloudImportedProvider[] = [];
for (const rawProvider of value.providers) {
const provider = parseCloudImportedProvider(rawProvider);
if (!provider) throw new Error("Invalid cloud provider sync provider response.");
providers.push(provider);
}
let lastRun: OpenworkCloudProviderSyncStatus["lastRun"] = null;
if ("lastRun" in value && value.lastRun !== null) {
if (!value.lastRun || typeof value.lastRun !== "object" || !("at" in value.lastRun) || (typeof value.lastRun.at !== "string" && typeof value.lastRun.at !== "number")) {
throw new Error("Invalid cloud provider sync last-run response.");
}
const run = parseCloudProviderSyncRun(value.lastRun);
lastRun = { at: value.lastRun.at, status: run.status, message: run.message };
}
return { hasSession: value.hasSession, lastRun, providers };
}

export type OpenworkServerStatus = "connected" | "disconnected" | "limited";

export type OpenworkServerDiagnostics = {
Expand Down Expand Up @@ -1334,6 +1398,24 @@ export function createOpenworkServerClient(options: { baseUrl: string; token?: s
const suffix = query.size ? `?${query.toString()}` : "";
return requestJson<OpenworkConnectState>(baseUrl, `/experimental/connect/state${suffix}`, { token, hostToken, timeoutMs: timeouts.config });
},
putDenSession: async (body: { baseUrl: string; token: string; orgId: string }) => {
await requestJson<unknown>(baseUrl, "/den-session", { hostToken, method: "PUT", body, timeoutMs: timeouts.config });
},
deleteDenSession: async () => {
await requestJson<unknown>(baseUrl, "/den-session", { hostToken, method: "DELETE", timeoutMs: timeouts.config });
},
runCloudProviderSyncNow: async (reason?: string) =>
parseCloudProviderSyncRun(await requestJson<unknown>(baseUrl, "/cloud-provider-sync/run", {
hostToken,
method: "POST",
body: reason ? { reason } : {},
timeoutMs: timeouts.cloudMcpReconcile,
})),
getCloudProviderSyncStatus: async () =>
parseCloudProviderSyncStatus(await requestJson<unknown>(baseUrl, "/cloud-provider-sync/status", {
token,
timeoutMs: timeouts.config,
})),
setConnectState: (connectEnabled: boolean) => requestJson<OpenworkConnectState>(baseUrl, "/experimental/connect/state", { token, hostToken, method: "PUT", body: { connectEnabled }, timeoutMs: timeouts.config }),
callExtensionAction: (payload: OpenworkExtensionActionCall) =>
requestJson<OpenworkExtensionActionResult>(baseUrl, "/experimental/extensions/call", {
Expand Down
96 changes: 94 additions & 2 deletions apps/app/src/react-app/domains/connections/provider-auth/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import { t } from "../../../../i18n";
import {
createDenClient,
readDenSettings,
resolveDenBaseUrls,
type DenOrgLlmProvider,
type DenOrgLlmProviderConnection,
} from "../../../../app/lib/den";
Expand Down Expand Up @@ -50,7 +51,8 @@ export type ProviderAuthOpenworkServer = {
OpenworkServerStoreSnapshot,
"openworkServerStatus" | "openworkServerClient"
> & {
openworkServerCapabilities: { config?: { read?: boolean; write?: boolean } } | null;
openworkServerAuth?: { token?: string; hostToken?: string };
openworkServerCapabilities: { config?: { read?: boolean; write?: boolean }; providerSync?: boolean } | null;
};
};
import {
Expand Down Expand Up @@ -268,6 +270,9 @@ export function createProviderAuthStore(options: CreateProviderAuthStoreOptions)
let cloudOrgProvidersInFlight: Promise<DenOrgLlmProvider[]> | null = null;
let cloudProviderSyncTail: Promise<void> = Promise.resolve();
let cloudProviderSyncContextKey = "";
let lastDenSessionPushKey = "";
let denSessionPushKey = "";
let denSessionPushInFlight: Promise<void> | null = null;

const emitChange = () => {
for (const listener of listeners) listener();
Expand Down Expand Up @@ -342,6 +347,42 @@ export function createProviderAuthStore(options: CreateProviderAuthStoreOptions)
};
};

const serverHandlesProviderSync = () => {
const openworkSnapshot = options.openworkServer.getSnapshot();
return Boolean(
openworkSnapshot.openworkServerStatus === "connected" &&
openworkSnapshot.openworkServerCapabilities?.providerSync === true &&
openworkSnapshot.openworkServerAuth?.hostToken?.trim() &&
openworkSnapshot.openworkServerClient,
);
};

const pushDenSession = (force = false): Promise<void> => {
const openworkSnapshot = options.openworkServer.getSnapshot();
const openworkClient = openworkSnapshot.openworkServerClient;
const settings = readDenSettings();
const apiBaseUrl = settings.apiBaseUrl ?? resolveDenBaseUrls(settings).apiBaseUrl;
const token = settings.authToken?.trim() ?? "";
const orgId = settings.activeOrgId?.trim() ?? "";
if (!serverHandlesProviderSync() || !openworkClient || !token || !orgId) return Promise.resolve();
const key = `${apiBaseUrl}::${orgId}::${token}`;
if (!force && key === lastDenSessionPushKey) return Promise.resolve();
if (key === denSessionPushKey && denSessionPushInFlight) return denSessionPushInFlight;
denSessionPushKey = key;
const request = openworkClient.putDenSession({ baseUrl: apiBaseUrl, token, orgId });
denSessionPushInFlight = request;
request.then(
() => { lastDenSessionPushKey = key; },
() => undefined,
).finally(() => {
if (denSessionPushInFlight === request) {
denSessionPushInFlight = null;
denSessionPushKey = "";
}
});
return request;
};

const refreshSnapshot = () => {
snapshot = {
providerAuthModalOpen: state.providerAuthModalOpen,
Expand Down Expand Up @@ -465,6 +506,14 @@ export function createProviderAuthStore(options: CreateProviderAuthStoreOptions)

const refreshImportedCloudProviders = async (refreshOptions?: { strict?: boolean }) => {
try {
if (serverHandlesProviderSync()) {
const openworkClient = options.openworkServer.getSnapshot().openworkServerClient;
if (!openworkClient) throw new Error("OpenWork server unavailable.");
const status = await openworkClient.getCloudProviderSyncStatus();
const next = Object.fromEntries(status.providers.map((provider) => [provider.cloudProviderId, provider]));
setStateField("importedCloudProviders", next);
return next;
}
const config = await readWorkspaceOpenworkConfigRecord();
const cloudImports = readWorkspaceCloudImports(config);
const next = cloudImports.providers;
Expand Down Expand Up @@ -1902,6 +1951,38 @@ export function createProviderAuthStore(options: CreateProviderAuthStoreOptions)
return { outcome: "handled_server_side" };
}

if (serverHandlesProviderSync()) {
try {
const openworkClient = options.openworkServer.getSnapshot().openworkServerClient;
if (!openworkClient) throw new Error("OpenWork server unavailable.");
let result = await openworkClient.runCloudProviderSyncNow(reason);
if (result.status === "no_session") {
await pushDenSession(true);
result = await openworkClient.runCloudProviderSyncNow(reason);
}
if (result.status === "failed" || result.status === "no_session") {
const message = logCloudProviderSyncError(
reason,
new Error(result.message ?? "Cloud provider sync failed."),
);
if (reason === "settings_cloud_opened") {
setStateField("providerAuthError", message);
}
return;
}
if (result.status === "applied") {
await refreshProviders({ force: true });
}
return { outcome: "handled_server_side" };
} catch (error) {
const message = logCloudProviderSyncError(reason, error);
if (reason === "settings_cloud_opened") {
setStateField("providerAuthError", message);
}
return;
}
}

const request = cloudProviderSyncTail
.catch(() => undefined)
.then(() =>
Expand Down Expand Up @@ -2059,6 +2140,13 @@ export function createProviderAuthStore(options: CreateProviderAuthStoreOptions)
setStateField("lastSyncError", {});
void refreshImportedCloudProviders();
}
if (serverHandlesProviderSync()) {
const nextSyncContextKey = getCloudProviderSyncContextKey();
if (nextSyncContextKey === cloudProviderSyncContextKey) return;
cloudProviderSyncContextKey = nextSyncContextKey;
void pushDenSession().then(() => runCloudProviderSync("app_launch"));
return;
}
if (!hasCloudProviderSyncPrerequisites()) {
cloudProviderSyncContextKey = "";
return;
Expand Down Expand Up @@ -2093,8 +2181,12 @@ export function createProviderAuthStore(options: CreateProviderAuthStoreOptions)
providerAuthMethods: {},
lastSyncError: {},
}));
void runCloudProviderSync("sign_in");
void pushDenSession().then(() => runCloudProviderSync("sign_in"));
} else {
if (serverHandlesProviderSync()) {
lastDenSessionPushKey = "";
void options.openworkServer.getSnapshot().openworkServerClient?.deleteDenSession().catch(() => undefined);
}
// Sign-out or error: remove all cloud-imported providers from the workspace
// Capture the full import records BEFORE clearing state
const importedProviders = { ...state.importedCloudProviders };
Expand Down
Loading
Loading