Skip to content

Commit 6a72454

Browse files
fix(server): a secretRef is handed out once even under concurrent use
consume read the stored value and deleted it in two store calls, so two concurrent calls with one ref could both read it. Consumption is now serialized. The client's single-flight key for answering a card encodes its ids structurally, so ids containing a colon can't collide. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
1 parent 13a6c31 commit 6a72454

3 files changed

Lines changed: 29 additions & 2 deletions

File tree

‎apps/server/src/secrets/SecretRequests.test.ts‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,8 @@ const withService = <A, E>(
4949
Layer.succeed(
5050
ServerSecretStore.ServerSecretStore,
5151
ServerSecretStore.ServerSecretStore.of({
52-
get: (name) => Effect.succeed(Option.fromNullishOr(stored.get(name))),
52+
// Yields like a real file read, so concurrent callers can interleave.
53+
get: (name) => Effect.yieldNow.pipe(Effect.as(Option.fromNullishOr(stored.get(name)))),
5354
set: (name, value) => Effect.sync(() => void stored.set(name, value)),
5455
create: (name, value) =>
5556
stored.has(name)
@@ -148,6 +149,26 @@ it.effect("a saved answer becomes a one-use ref, and the thread only learns it w
148149
),
149150
);
150151

152+
it.effect("two concurrent uses of one ref hand the value out once", () =>
153+
withService(({ service }) =>
154+
Effect.gen(function* () {
155+
yield* service.answer({
156+
threadId,
157+
turnItemId,
158+
answer: { type: "save", secret: "ghp_secret" },
159+
});
160+
const ref = Option.getOrThrow(yield* service.savedRef({ threadId, turnItemId }));
161+
const results = yield* Effect.all(
162+
[service.consume({ ref, projectId }), service.consume({ ref, projectId })].map(
163+
Effect.result,
164+
),
165+
{ concurrency: "unbounded" },
166+
);
167+
assert.deepEqual(results.map((result) => result._tag).toSorted(), ["Failure", "Success"]);
168+
}),
169+
),
170+
);
171+
151172
it.effect("a ref only works in the project it was entered for", () =>
152173
withService(({ service }) =>
153174
Effect.gen(function* () {

‎apps/server/src/secrets/SecretRequests.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import * as Layer from "effect/Layer";
2727
import * as Option from "effect/Option";
2828
import * as Schedule from "effect/Schedule";
2929
import * as Schema from "effect/Schema";
30+
import * as Semaphore from "effect/Semaphore";
3031

3132
import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
3233
import * as Metrics from "../observability/Metrics.ts";
@@ -178,8 +179,12 @@ const make = Effect.gen(function* () {
178179
),
179180
);
180181

182+
// get and remove are separate store calls; one consumer at a time keeps two
183+
// concurrent calls from both reading a ref before either deletes it.
184+
const consumeLock = yield* Semaphore.make(1);
181185
const consume: SecretRequests["Service"]["consume"] = (input) =>
182186
consumeRef(input).pipe(
187+
consumeLock.withPermits(1),
183188
Effect.tap(() => Metrics.increment(Metrics.secretRefsConsumedTotal, { result: "used" })),
184189
Effect.tapError(() =>
185190
Metrics.increment(Metrics.secretRefsConsumedTotal, { result: "rejected" }),

‎packages/client-runtime/src/state/server.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1306,7 +1306,8 @@ export function createServerEnvironmentAtoms<R, E>(
13061306
tag: WS_METHODS.secretsAnswerRequest,
13071307
concurrency: {
13081308
mode: "singleFlight",
1309-
key: ({ environmentId, input }) => `${environmentId}:${input.threadId}:${input.turnItemId}`,
1309+
key: ({ environmentId, input }) =>
1310+
JSON.stringify([environmentId, input.threadId, input.turnItemId]),
13101311
},
13111312
}),
13121313
refreshUsageRates: createEnvironmentRpcCommand(runtime, {

0 commit comments

Comments
 (0)