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
57 changes: 53 additions & 4 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,7 @@ export const OpenCodeMemPlugin: Plugin = async (ctx: PluginInput) => {
// Aborted on plugin dispose: queued-but-not-started auto-capture jobs are
// skipped; in-flight work keeps its original (uncancellable) semantics.
const pluginLifetime = new AbortController();
let cleanedUp = false;

const GLOBAL_PLUGIN_WARMUP_KEY = Symbol.for("opencode-mem.plugin.warmedup");

Expand All @@ -328,7 +329,11 @@ export const OpenCodeMemPlugin: Plugin = async (ctx: PluginInput) => {

await configureOpencodeHostTransport(ctx);

(async () => {
// Re-reads the provider directory into the connectivity snapshot. Errors
// are logged and the previous snapshot retained. The host registers
// user-config providers after plugin setup (provider.updated/model.updated),
// so this runs again on those events and on connectivity-gate misses.
const refreshConnectedProviders = async (): Promise<void> => {
try {
const providerResult = await ctx.client.provider.list();
if (providerResult.data?.connected) {
Expand All @@ -344,9 +349,46 @@ export const OpenCodeMemPlugin: Plugin = async (ctx: PluginInput) => {
});
}
} catch (error) {
log("Failed to initialize opencode provider state", { error: String(error) });
log("Failed to refresh opencode provider state", { error: String(error) });
}
})();
};

// Coalesce refresh triggers into one in-flight pass plus trailing drain.
// Event handlers and ensureProviderConnected (gate miss) share this path so
// concurrent list() calls collapse; callers can await the same promise.
let refreshInFlight: Promise<void> | undefined;
let refreshQueued = false;
const runCoalescedRefresh = (): Promise<void> => {
refreshQueued = true;
if (refreshInFlight) return refreshInFlight;

refreshInFlight = (async () => {
while (refreshQueued && !cleanedUp) {
refreshQueued = false;
await refreshConnectedProviders();
}
})().finally(() => {
refreshInFlight = undefined;
if (refreshQueued && !cleanedUp) {
void runCoalescedRefresh();
}
});
return refreshInFlight;
};

void runCoalescedRefresh();

// Gate miss awaits the same coalesced refresh; cleared on dispose so no
// refresh fires afterwards. Only clear our own refresher identity.
const { setConnectedProvidersRefresher, clearConnectedProvidersRefresher } =
await loadOpencodeProvider();
const providerRefresher = async () => {
await runCoalescedRefresh();
};
setConnectedProvidersRefresher(providerRefresher);
const clearProviderRefreshWiring = () => {
clearConnectedProvidersRefresher(providerRefresher);
};

let tursoReadyForWeb = !isConfigured();
if (CONFIG.webServerEnabled && isConfigured()) {
Expand Down Expand Up @@ -471,10 +513,10 @@ export const OpenCodeMemPlugin: Plugin = async (ctx: PluginInput) => {
});
}

let cleanedUp = false;
const cleanupPlugin = async () => {
if (cleanedUp) return;
cleanedUp = true;
clearProviderRefreshWiring();
pluginLifetime.abort();
for (const timer of idleTimers.values()) clearTimeout(timer);
idleTimers.clear();
Expand Down Expand Up @@ -1076,6 +1118,13 @@ export const OpenCodeMemPlugin: Plugin = async (ctx: PluginInput) => {
event: async (input: { event: { type: string; properties?: any } }) => {
const event = input.event;

// Late host provider registration: user-config providers can appear
// (or disappear) after plugin setup; refresh the connectivity snapshot.
if (event.type === "provider.updated" || event.type === "model.updated") {
if (!cleanedUp) runCoalescedRefresh();
return;
}

// Client-side step watchdog for internal structured-output sessions (#278).
// OpenCode's agent.steps soft-cap does not hard-stop json_schema loops when
// forced StructuredOutput keeps failing (e.g. opencode-claude-auth).
Expand Down
31 changes: 31 additions & 0 deletions src/services/ai/opencode-provider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,9 @@ let _v2Client: OpencodeClient | undefined;
let _v2BaseUrl: string | undefined;
let _hostFetch: typeof fetch | undefined;
let _useSdkTransport = false;
// Registered by the plugin init so a gate miss can force one provider-directory
// refresh before failing (late host provider registration).
let _connectedProvidersRefresher: (() => Promise<void>) | undefined;

export function setHostFetch(customFetch: typeof fetch): void {
_hostFetch = customFetch;
Expand All @@ -105,6 +108,34 @@ export function isProviderConnected(providerName: string): boolean {
return _connectedProviders.has(providerName);
}

/** Register (or clear) the refresh callback used by ensureProviderConnected. */
export function setConnectedProvidersRefresher(refresher?: () => Promise<void>): void {
_connectedProvidersRefresher = refresher;
}

/** Clear the refresh callback only if it is still the one being disposed. */
export function clearConnectedProvidersRefresher(refresher: () => Promise<void>): void {
if (_connectedProvidersRefresher === refresher) _connectedProvidersRefresher = undefined;
}

/**
* Refresh-on-miss for the connectivity gate: when a provider id is missing
* from the snapshot and a refresher is registered, await one refresh and
* re-check. Returns the re-checked connectivity; never refreshes twice for
* the same call and never throws (refresh errors are logged by the refresher).
*/
export async function ensureProviderConnected(providerName: string): Promise<boolean> {
if (_connectedProviders.has(providerName)) return true;
const refresher = _connectedProvidersRefresher;
if (!refresher) return false;
try {
await refresher();
} catch {
// The refresher logs; a failed refresh keeps the previous set.
}
return _connectedProviders.has(providerName);
}

export function setV2Client(client: OpencodeClient): void {
_v2Client = client;
// Native v2 adapters pass a session-capable client without a server URL.
Expand Down
6 changes: 4 additions & 2 deletions src/services/ai/profile-llm-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,11 @@ export async function getOpenCodeClient(): Promise<OpencodeClient> {
return _cachedClient;
}

const { isProviderConnected, getV2Client } = await loadOpencodeProvider();
const { ensureProviderConnected, getV2Client } = await loadOpencodeProvider();

if (!isProviderConnected(provider)) {
// One refresh-on-miss before failing: user-config providers register
// seconds after plugin setup, so the init snapshot may lag the host.
if (!(await ensureProviderConnected(provider))) {
throw new Error(
`opencode provider '${provider}' is not connected. Check your opencode provider configuration.`
);
Expand Down
6 changes: 4 additions & 2 deletions src/services/auto-capture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -529,7 +529,7 @@ async function generateSummary(
log("opencodeProvider takes precedence over memoryModel for auto-capture");
}

const { isProviderConnected, getV2Client, generateStructuredOutput } =
const { ensureProviderConnected, getV2Client, generateStructuredOutput } =
await loadOpencodeProvider();

// "inherit" resolves to the model opencode used for the captured prompt
Expand All @@ -547,7 +547,9 @@ async function generateSummary(
modelID = prompt.modelId;
}

if (!isProviderConnected(providerID)) {
// One refresh-on-miss before failing: user-config providers register
// seconds after plugin setup, so the init snapshot may lag the host.
if (!(await ensureProviderConnected(providerID))) {
throw new Error(
`opencode provider '${providerID}' is not connected. Check your opencode provider configuration.`
);
Expand Down
31 changes: 20 additions & 11 deletions src/v2/legacy-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -283,16 +283,25 @@ export async function eventBelongsToLocation(ctx: Context, raw: any): Promise<bo

const data = raw?.data ?? raw?.payload?.data ?? raw?.properties;
const sessionID = data?.sessionID ?? data?.session?.id ?? data?.info?.id;
if (!sessionID) return false;
try {
const session: any = await ctx.session.get({ sessionID });
const sessionDirectory =
session?.location?.directory ?? session?.directory ?? session?.data?.directory;
return (
typeof sessionDirectory === "string" &&
resolve(sessionDirectory) === resolve(ctx.location.directory)
);
} catch {
return false;
if (sessionID) {
try {
const session: any = await ctx.session.get({ sessionID });
const sessionDirectory =
session?.location?.directory ?? session?.directory ?? session?.data?.directory;
return (
typeof sessionDirectory === "string" &&
resolve(sessionDirectory) === resolve(ctx.location.directory)
);
} catch {
return false;
}
}

// Provider/model inventory updates are process-scoped ephemeral events with
// optional location. When the host omits location (and there is no session),
// still forward them so connectivity snapshots can refresh on late registration.
const envelope = raw?.payload ?? raw;
const source = envelope?.type === "sync" && envelope.syncEvent ? envelope.syncEvent : envelope;
const rawType = typeof source?.type === "string" ? source.type.replace(/\.1$/, "") : undefined;
return rawType === "provider.updated" || rawType === "model.updated";
}
8 changes: 4 additions & 4 deletions tests/auto-capture-concurrency.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -311,7 +311,7 @@ const llmEntered = new Promise((res) => { llmEnteredResolve = res; });
const llmOrder = [];
mock.module(${JSON.stringify(u("src/services/ai/opencode-provider-loader.js"))}, () => ({
loadOpencodeProvider: async () => ({
isProviderConnected: () => true,
ensureProviderConnected: async () => true,
getV2Client: () => ({}),
generateStructuredOutput: async ({ userPrompt }) => {
const sid = userPrompt.includes("Fix login bug") ? "sess-A" : "sess-B";
Expand Down Expand Up @@ -425,7 +425,7 @@ mock.module(${JSON.stringify(u("src/services/tags.js"))}, () => ({
}));
mock.module(${JSON.stringify(u("src/services/ai/opencode-provider-loader.js"))}, () => ({
loadOpencodeProvider: async () => ({
isProviderConnected: () => true,
ensureProviderConnected: async () => true,
getV2Client: () => ({}),
generateStructuredOutput: async ({ userPrompt }) => {
const sid = userPrompt.includes("Do work") ? "sess-X" : "?";
Expand Down Expand Up @@ -501,7 +501,7 @@ mock.module(${JSON.stringify(u("src/services/tags.js"))}, () => ({
}));
mock.module(${JSON.stringify(u("src/services/ai/opencode-provider-loader.js"))}, () => ({
loadOpencodeProvider: async () => ({
isProviderConnected: () => true,
ensureProviderConnected: async () => true,
getV2Client: () => ({}),
generateStructuredOutput: async ({ userPrompt }) => {
if (userPrompt.includes("Will fail")) throw new Error("llm exploded");
Expand Down Expand Up @@ -607,7 +607,7 @@ let releaseA;
const gateA = new Promise((res) => { releaseA = res; });
mock.module(${JSON.stringify(u("src/services/ai/opencode-provider-loader.js"))}, () => ({
loadOpencodeProvider: async () => ({
isProviderConnected: () => true,
ensureProviderConnected: async () => true,
getV2Client: () => ({}),
generateStructuredOutput: async ({ userPrompt }) => {
const sid = userPrompt.includes("In flight") ? "sess-A" : "sess-B";
Expand Down
4 changes: 2 additions & 2 deletions tests/auto-capture.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -157,7 +157,7 @@ mock.module(${JSON.stringify(languageUrl)}, () => ({
}));
mock.module(${JSON.stringify(opencodeProviderLoaderUrl)}, () => ({
loadOpencodeProvider: async () => ({
isProviderConnected: () => true,
ensureProviderConnected: async () => true,
getV2Client: () => ({}),
generateStructuredOutput: async ({ userPrompt }) => {
summaryPrompts.push(userPrompt);
Expand Down Expand Up @@ -286,7 +286,7 @@ mock.module(${JSON.stringify(languageUrl)}, () => ({
}));
mock.module(${JSON.stringify(opencodeProviderLoaderUrl)}, () => ({
loadOpencodeProvider: async () => ({
isProviderConnected: () => true,
ensureProviderConnected: async () => true,
getV2Client: () => ({}),
generateStructuredOutput: async () => {
throw new Error(
Expand Down
Loading
Loading