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
18 changes: 18 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,12 @@ rotations during legacy migration, and hardens file snapshots and operator
recovery. Artifact quotas now reserve 8 KiB for frame metadata, reducing the
available content allowance by 4 KiB, and apply to historical snapshots too.

A third pass prevents dispatch after cancellation, rejects malformed executor
replies, and closes connector scopes after human auth operations. OAuth handoffs
claim ownership atomically when storage supports it. Activity reads now recover
from malformed pages and ignore superseded pagination, and memory and D1 storage
reject invalid expiry values.

### Changed

- Pin Effect to `4.0.0`, replacing `4.0.0-rc.117`, and align the optional
Expand All @@ -25,6 +31,18 @@ available content allowance by 4 KiB, and apply to historical snapshots too.

### Fixed

- Refuse already-aborted deadline operations and preserve explicit `null`
cancellation reasons (#660) (#661).
- Return structured failures for malformed executor replies (#662).
- Close callback and credential-test connector scopes after their final use,
keeping cleanup bounded and best-effort (#663) (#664).
- Prevent concurrent OAuth state handoffs from replacing another owner on
storage that supports atomic claims (#665).
- Show retryable Activity errors for unreadable pages and discard superseded
pagination responses (#666) (#667).
- Reject invalid TTLs in memory and D1 storage before mutation, and match NUL
bytes literally in D1 list prefixes (#668) (#669).

- Enforce catalog count and byte ceilings even when discovery is invalidated
while the listing is in flight (#644).
- Retain OAuth cleanup lineage after uncertain storage commits, preserve
Expand Down
32 changes: 27 additions & 5 deletions examples/worker/src/d1-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,17 @@ import type { KVStorage } from "@zackbart/connecta";
export function d1Storage(db: D1Database): KVStorage {
// Every statement that tests liveness binds the current time as ?2.
const live = "(expires_at_ms IS NULL OR expires_at_ms > ?2)";
const expiresAt = (now: number, ttlSeconds?: number) =>
ttlSeconds ? now + ttlSeconds * 1000 : null;
const expiresAt = (now: number, ttlSeconds?: number): number | null => {
if (ttlSeconds !== undefined && !Number.isFinite(ttlSeconds)) {
throw new RangeError("ttlSeconds must produce a finite expiration timestamp");
}
if (!ttlSeconds) return null;
const expiry = now + ttlSeconds * 1000;
if (!Number.isFinite(expiry)) {
throw new RangeError("ttlSeconds must produce a finite expiration timestamp");
}
return expiry;
};
const changed = (result: D1Result) => result.meta.changes > 0;
return {
async get(key) {
Expand All @@ -31,6 +40,7 @@ export function d1Storage(db: D1Database): KVStorage {
},
async set(key, value, opts) {
const now = Date.now();
const expiry = expiresAt(now, opts?.ttlSeconds);
await db.batch([
db
.prepare(
Expand All @@ -46,7 +56,7 @@ export function d1Storage(db: D1Database): KVStorage {
ON CONFLICT (key) DO UPDATE SET
value = excluded.value, expires_at_ms = excluded.expires_at_ms`,
)
.bind(key, value, expiresAt(now, opts?.ttlSeconds)),
.bind(key, value, expiry),
]);
},
async delete(key) {
Expand All @@ -56,7 +66,7 @@ export function d1Storage(db: D1Database): KVStorage {
const { results } = await db
.prepare(
`SELECT key FROM connecta_kv
WHERE key >= ?1 AND substr(key, 1, length(?1)) = ?1 AND ${live}`,
WHERE key >= ?1 AND substr(CAST(key AS BLOB), 1, length(CAST(?1 AS BLOB))) = CAST(?1 AS BLOB) AND ${live}`,
)
.bind(prefix, Date.now())
.all<{ key: string }>();
Expand Down Expand Up @@ -84,7 +94,19 @@ export function d1Storage(db: D1Database): KVStorage {
.run(),
);
}
const expiry = expiresAt(now, opts?.ttlSeconds);
let expiry: number | null;
try {
expiry = expiresAt(now, opts?.ttlSeconds);
} catch (error) {
// Invalid write options do not change a failed comparison's result.
// Ordinary claims still use one atomic statement without this read.
const current = await db
.prepare(`SELECT value FROM connecta_kv WHERE key = ?1 AND ${live}`)
.bind(key, now)
.first<{ value: string }>();
if ((current?.value ?? null) !== expected) return false;
throw error;
}
if (expected === null) {
// Insert, or take over a row that has expired. A live row makes the
// upsert's WHERE false, so nothing changes and the claim is refused.
Expand Down
31 changes: 31 additions & 0 deletions src/connector-scope.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,37 @@ import type { Connector, ConnectorContext } from "./types.js";
import { closeScope } from "./runtime/connector-scope.js";
import { runEdge } from "./runtime/run.js";

// Internal ownership for credential probes invoked by a route. Context copies
// keep the requestScope, and each connector retains its own cleanup ownership.
// Standalone connector tests still close themselves.
const cleanupOwners = new WeakMap<object, Map<string, number>>();

export function claimConnectorScopeCleanup(
ctx: ConnectorContext,
connectorId: string,
): () => void {
const scope = ctx.requestScope ?? ctx;
let owners = cleanupOwners.get(scope);
if (!owners) cleanupOwners.set(scope, owners = new Map());
owners.set(connectorId, (owners.get(connectorId) ?? 0) + 1);
let released = false;
return () => {
if (released) return;
released = true;
const remaining = owners.get(connectorId)! - 1;
if (remaining) owners.set(connectorId, remaining);
else owners.delete(connectorId);
if (!owners.size) cleanupOwners.delete(scope);
};
}

export function connectorScopeCleanupClaimed(
ctx: ConnectorContext,
connectorId: string,
): boolean {
return cleanupOwners.get(ctx.requestScope ?? ctx)?.has(connectorId) === true;
}

/** Runtime hook for work that may safely continue after a response is ready. */
export type DeferredWork = (promise: Promise<unknown>) => void;

Expand Down
7 changes: 4 additions & 3 deletions src/connectors/remote-mcp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import {
} from "../auth/downstream-oauth.js";
import { trackOAuthStartReset } from "../auth/oauth-start-reset.js";
import { MAX_CATALOG_TOOLS } from "../catalog-limits.js";
import { connectorScopeCleanupClaimed } from "../connector-scope.js";
import {
boundedEchoText,
ConnectorCallError,
Expand Down Expand Up @@ -1402,9 +1403,9 @@ export function remoteMcp(id: string, opts: RemoteMcpOptions): Connector {
} catch (err) {
return { ok: false, message: msg(err) };
} finally {
// A test owns the scope it just opened; leaving the session for
// the downstream to age out is not this button's to spend.
await connector.closeScope?.(ctx);
// Standalone tests own their scope; an operator route owns
// bounded teardown for the context it supplied to this hook.
if (!connectorScopeCleanupClaimed(ctx, id)) await connector.closeScope?.(ctx);
}
},
}
Expand Down
15 changes: 15 additions & 0 deletions src/execute.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1447,7 +1447,22 @@ function finishedRun(
outcome: ExecuteResult,
{ emitted, diagnostics, invocationFailures }: RunReport,
): ToolResult {
if (outcome === null || typeof outcome !== "object" || Array.isArray(outcome)) {
return failureResponse("Executor failed: expected an ExecuteResult object.", {
emitted,
diagnostics,
code: "executor_failed",
});
}
const logs = executeLogs(outcome.logs);
if (outcome.error !== undefined && typeof outcome.error !== "string") {
return failureResponse("Executor failed: ExecuteResult.error must be a string.", {
logs,
emitted,
diagnostics,
code: "executor_failed",
});
}
if (outcome.error !== undefined) {
// Executor bridges necessarily reduce thrown host errors to strings.
// Match that terminal string back to the request-local typed failure so
Expand Down
6 changes: 5 additions & 1 deletion src/operator-ui/app/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -617,9 +617,12 @@ export function testCredential(connector: string): Promise<void> {

/* Activity ---------------------------------------------------------------- */

let activityRevision = 0;
export async function loadActivity(reset: boolean): Promise<void> {
if (!state.data?.activityEnabled) return;
const current = fence();
const identityCurrent = fence();
const revision = ++activityRevision;
const current = () => identityCurrent() && revision === activityRevision;
set({
activityPhase: "loading",
activityNotice: null,
Expand All @@ -636,6 +639,7 @@ export async function loadActivity(reset: boolean): Promise<void> {
current,
);
if (!current()) return;
if (!payload || !Array.isArray(payload.events)) throw new RequestFailure("refused");
set({
activityPhase: "ready",
activityEvents: [
Expand Down
2 changes: 1 addition & 1 deletion src/operator-ui/generated.ts

Large diffs are not rendered by default.

27 changes: 17 additions & 10 deletions src/registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -612,17 +612,24 @@ export class Registry implements RegistryView {
principalKey: string,
): Promise<void> {
const key = this.oauthHandoffKey(connectorId, await sha256Hex(state));
const existing = await this.opts.storage.get(key);
if (existing && existing !== principalKey) {
throw new Error(
`Connector "${connectorId}" reused one OAuth state across principals`,
);
for (let attempt = 0; attempt < 32; attempt++) {
const existing = await this.opts.storage.get(key);
if (existing && existing !== principalKey) {
throw new Error(
`Connector "${connectorId}" reused one OAuth state across principals`,
);
}
const expiry = { ttlSeconds: OAUTH_HANDOFF_TTL_SECONDS };
if (this.opts.storage.compareAndSet) {
// Bind ownership atomically, including renewal of the same owner's
// handoff. Two owners reading a miss must not overwrite each other.
if (!await this.opts.storage.compareAndSet(key, existing, principalKey, expiry)) continue;
} else {
await this.opts.storage.set(key, principalKey, expiry);
}
return;
}
await this.opts.storage.set(
key,
principalKey,
{ ttlSeconds: OAUTH_HANDOFF_TTL_SECONDS },
);
throw new Error(`OAuth handoff for "${connectorId}" is busy; retry authorization`);
}

async oauthCallbackView(
Expand Down
29 changes: 18 additions & 11 deletions src/routes/credentials.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { Effect } from "effect";
import { claimConnectorScopeCleanup, closeConnectorScope } from "../connector-scope.js";
import {
credentialTestRule,
describeCredentialTestMismatch,
Expand Down Expand Up @@ -210,17 +211,23 @@ function credentialRequest(
}
const storedValues = values!;
const ctx = registry.contextFor(connectorId, baseUrl);
const result =
mode === "multiple"
? await connector.testCredentials!(storedValues, ctx)
: await connector.testCredential!(
// The single-value shape check above guarantees this key.
storedValues.value!,
ctx,
);
const ok = result?.ok === true;
logged(ok, result?.message);
return privateJson({ ok });
const releaseCleanup = claimConnectorScopeCleanup(ctx, connectorId);
try {
const result =
mode === "multiple"
? await connector.testCredentials!(storedValues, ctx)
: await connector.testCredential!(
// The single-value shape check above guarantees this key.
storedValues.value!,
ctx,
);
const ok = result?.ok === true;
logged(ok, result?.message);
return privateJson({ ok });
} finally {
releaseCleanup();
await closeConnectorScope(connector, ctx, context.defer);
}
},
catch: (error) => error,
}).pipe(
Expand Down
Loading
Loading