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
10 changes: 3 additions & 7 deletions apps/server/src/orchestration-v2/Adapters/OpenCode2AdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ import * as McpProviderSession from "../../mcp/McpProviderSession.ts";
import { buildRuntimeInstructions } from "../../provider/RuntimeInstructions.ts";
import { t3OrchestrationSystemPrompt } from "../../provider/T3OrchestrationInstructions.ts";
import { SKILL_MENTION_PATTERN } from "@t3tools/shared/composerInlineTokens";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import { getModelSelectionStringOptionValue, modelSelectionsEqual } from "@t3tools/shared/model";
import { causeErrorTag } from "@t3tools/shared/observability";

Expand Down Expand Up @@ -887,18 +888,13 @@ export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: Pr
const stagedReverts = new Set<string>();
// Starting a turn and cutting the history take turns on a session: each
// checks that the other is not running before its own requests yield.
const sessionGates = new Map<string, Semaphore.Semaphore>();
const sessionGates = yield* KeyedLock.make<string>();
const exclusive =
(providerThread: OrchestrationV2ProviderThread) =>
<A, E, R>(effect: Effect.Effect<A, E, R>) => {
const sessionId = providerThread.nativeThreadRef?.nativeId;
if (sessionId == null) return effect;
let gate = sessionGates.get(sessionId);
if (gate === undefined) {
gate = Semaphore.makeUnsafe(1);
sessionGates.set(sessionId, gate);
}
return gate.withPermit(effect);
return sessionGates.withLock(sessionId, effect);
};
const emit = (event: ProviderAdapter.ProviderAdapterV2Event) =>
Queue.offer(events, event).pipe(Effect.asVoid);
Expand Down
27 changes: 3 additions & 24 deletions apps/server/src/orchestration-v2/CheckpointService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,10 +14,9 @@ import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Encoding from "effect/Encoding";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import * as Layer from "effect/Layer";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as Semaphore from "effect/Semaphore";

import { parseTurnDiffFilesFromNumstat } from "../checkpointing/Diffs.ts";
import * as CheckpointStore from "../checkpointing/CheckpointStore.ts";
Expand Down Expand Up @@ -249,29 +248,9 @@ export const layer: Layer.Layer<
Effect.gen(function* () {
const checkpointStore = yield* CheckpointStore.CheckpointStore;
const idAllocator = yield* IdAllocator.IdAllocatorV2;
const workspaceSemaphores = yield* Ref.make(new Map<string, Semaphore.Semaphore>());

const getWorkspaceSemaphore = (cwd: string) =>
Effect.gen(function* () {
const existing = (yield* Ref.get(workspaceSemaphores)).get(cwd);
if (existing !== undefined) {
return existing;
}

const created = yield* Semaphore.make(1);
return yield* Ref.modify(workspaceSemaphores, (current) => {
const concurrent = current.get(cwd);
if (concurrent !== undefined) {
return [concurrent, current];
}
const updated = new Map(current);
updated.set(cwd, created);
return [created, updated];
});
});

const workspaceLocks = yield* KeyedLock.make<string>();
const withWorkspaceLock = <A, E, R>(cwd: string, effect: Effect.Effect<A, E, R>) =>
Effect.flatMap(getWorkspaceSemaphore(cwd), (semaphore) => semaphore.withPermits(1)(effect));
workspaceLocks.withLock(cwd, effect);

const isGitCheckpointable = (cwd: string) =>
checkpointStore.isGitRepository(cwd).pipe(Effect.orElseSucceed(() => false));
Expand Down
67 changes: 0 additions & 67 deletions apps/server/src/orchestration-v2/KeyedSerialExecutor.test.ts

This file was deleted.

8 changes: 4 additions & 4 deletions apps/server/src/orchestration-v2/ProviderSessionManager.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import { resolveProjectSettings } from "@t3tools/shared/projectSettings";
import {
ModelSelection,
Expand Down Expand Up @@ -36,7 +37,6 @@ import * as ServerSettings from "../serverSettings.ts";
import * as McpSessionRegistry from "../mcp/McpSessionRegistry.ts";
import * as EventSink from "./EventSink.ts";
import * as IdAllocator from "./IdAllocator.ts";
import { makeKeyedSerialExecutor } from "./KeyedSerialExecutor.ts";
import * as ProviderEventIngestor from "./ProviderEventIngestor.ts";
import {
ProviderAdapterEventStreamError,
Expand Down Expand Up @@ -386,9 +386,9 @@ export const layerWithOptions = (
// cannot drop cleanup for threads only the earlier session served.
const releaseRecordRetries = yield* FiberSet.make();
const nextSubscriberId = yield* Ref.make(0);
const sessionOpen = yield* makeKeyedSerialExecutor<ProviderSessionId>();
const sessionOpen = yield* KeyedLock.make<ProviderSessionId>();
// Orders a thread's attach against a detach unloading it on the same session.
const threadAttachment = yield* makeKeyedSerialExecutor<string>();
const threadAttachment = yield* KeyedLock.make<string>();
const threadAttachmentKey = (input: {
readonly providerSessionId: ProviderSessionId;
readonly threadId: ThreadId;
Expand Down Expand Up @@ -427,7 +427,7 @@ export const layerWithOptions = (
};
const isMcpCredentialReserved = (threadId: ThreadId, mcpCredentialId: string) =>
(mcpCredentialReservations.get(mcpReservationKey(threadId, mcpCredentialId)) ?? 0) > 0;
const mcpPrepareLock = yield* makeKeyedSerialExecutor<ThreadId>();
const mcpPrepareLock = yield* KeyedLock.make<ThreadId>();
/**
* Resolves (or mints) the thread's MCP credential and returns it with a
* reservation held; the caller must drop the reservation exactly once.
Expand Down
7 changes: 3 additions & 4 deletions apps/server/src/orchestration-v2/ThreadCommandExecutor.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,12 @@
import type { ThreadId } from "@t3tools/contracts";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import * as Context from "effect/Context";
import * as Layer from "effect/Layer";

import { makeKeyedSerialExecutor, type KeyedSerialExecutor } from "./KeyedSerialExecutor.ts";

/** Shared by thread commands and project deletion so both plan against current thread state. */
export class ThreadCommandExecutor extends Context.Service<
ThreadCommandExecutor,
KeyedSerialExecutor<ThreadId>
KeyedLock.KeyedLock<ThreadId>
>()("t3/orchestration-v2/ThreadCommandExecutor") {}

export const layer = Layer.effect(ThreadCommandExecutor, makeKeyedSerialExecutor<ThreadId>());
export const layer = Layer.effect(ThreadCommandExecutor, KeyedLock.make<ThreadId>());
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import {
ThreadPullRequestLink,
TurnItemId,
} from "@t3tools/contracts";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
Expand All @@ -30,7 +31,6 @@ import * as Schema from "effect/Schema";
import * as SqlClient from "effect/unstable/sql/SqlClient";

import * as EventSink from "../EventSink.ts";
import { makeKeyedSerialExecutor } from "../KeyedSerialExecutor.ts";
import { randomUuidV4 } from "../RandomUuid.ts";

const IMPORT_EVENT_PREFIX = "migration:v1";
Expand Down Expand Up @@ -347,7 +347,7 @@ function chunks<A>(items: ReadonlyArray<A>, size: number): Array<ReadonlyArray<A
const make = Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const eventSink = yield* EventSink.EventSinkV2;
const transcriptImports = yield* makeKeyedSerialExecutor<ThreadId>();
const transcriptImports = yield* KeyedLock.make<ThreadId>();

const listMessages = (threadId: ThreadId) =>
sql<LegacyMessageRow>`
Expand Down
6 changes: 3 additions & 3 deletions apps/server/src/project/ProjectService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
type ProjectSnapshot,
type ThreadId,
} from "@t3tools/contracts";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
Expand All @@ -18,7 +19,6 @@ import * as Schema from "effect/Schema";

import * as EventSink from "../orchestration-v2/EventSink.ts";
import * as IdAllocator from "../orchestration-v2/IdAllocator.ts";
import { makeKeyedSerialExecutor } from "../orchestration-v2/KeyedSerialExecutor.ts";
import * as LegacyV1ThreadImporter from "../orchestration-v2/legacy/LegacyV1ThreadImporter.ts";
import * as ProjectionStore from "../orchestration-v2/ProjectionStore.ts";
import {
Expand Down Expand Up @@ -155,8 +155,8 @@ export const make = Effect.gen(function* () {
const threadCommands = yield* ThreadCommandExecutor.ThreadCommandExecutor;
// Commands for one project run in order. Commands that claim a workspace root
// also hold that root, so two projects cannot both claim it.
const projectLocks = yield* makeKeyedSerialExecutor<ProjectId>();
const workspaceLocks = yield* makeKeyedSerialExecutor<string>();
const projectLocks = yield* KeyedLock.make<ProjectId>();
const workspaceLocks = yield* KeyedLock.make<string>();

const toProject = (
row: ProjectStore.ProjectRow,
Expand Down
33 changes: 6 additions & 27 deletions apps/server/src/provider/providerMaintenanceCommandCoordinator.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import * as Effect from "effect/Effect";
import * as Ref from "effect/Ref";
import * as Semaphore from "effect/Semaphore";

export interface ProviderMaintenanceCommandCoordinatorShape<E> {
readonly withCommandLock: <A, R>(input: {
Expand All @@ -15,7 +15,7 @@ export const makeProviderMaintenanceCommandCoordinator = Effect.fn(
"makeProviderMaintenanceCommandCoordinator",
)(function* <E>(input: { readonly makeAlreadyRunningError: (targetKey: string) => E }) {
const runningTargetsRef = yield* Ref.make<ReadonlySet<string>>(new Set());
const locksRef = yield* Ref.make<ReadonlyMap<string, Semaphore.Semaphore>>(new Map());
const locks = yield* KeyedLock.make<string>();

const acquireTarget = Effect.fn("acquireTarget")(function* (targetKey: string) {
return yield* Ref.modify(runningTargetsRef, (runningTargets) => {
Expand All @@ -35,24 +35,6 @@ export const makeProviderMaintenanceCommandCoordinator = Effect.fn(
return next;
});

const getLock = Effect.fn("getProviderMaintenanceCommandLock")(function* (lockKey: string) {
const existing = (yield* Ref.get(locksRef)).get(lockKey);
if (existing) {
return existing;
}

const lock = yield* Semaphore.make(1);
return yield* Ref.modify(locksRef, (locks) => {
const current = locks.get(lockKey);
if (current) {
return [current, locks] as const;
}
const next = new Map(locks);
next.set(lockKey, lock);
return [lock, next] as const;
});
});

const withCommandLock: ProviderMaintenanceCommandCoordinatorShape<E>["withCommandLock"] = ({
targetKey,
lockKey,
Expand All @@ -65,13 +47,10 @@ export const makeProviderMaintenanceCommandCoordinator = Effect.fn(
return yield* Effect.fail(input.makeAlreadyRunningError(targetKey));
}

return yield* Effect.gen(function* () {
const lock = yield* getLock(lockKey);
if (onQueued) {
yield* onQueued;
}
return yield* lock.withPermits(1)(run);
}).pipe(Effect.ensuring(releaseTarget(targetKey)));
return yield* (onQueued ?? Effect.void).pipe(
Effect.andThen(locks.withLock(lockKey, run)),
Effect.ensuring(releaseTarget(targetKey)),
);
});

return {
Expand Down
25 changes: 3 additions & 22 deletions apps/server/src/terminal/Manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ import {
ProviderInstanceId,
} from "@t3tools/contracts";
import { makeKeyedCoalescingWorker } from "@t3tools/shared/KeyedCoalescingWorker";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess";
import { mergePathEntries } from "@t3tools/shared/shell";

Expand All @@ -58,7 +59,6 @@ import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Schema from "effect/Schema";
import * as Scope from "effect/Scope";
import * as Semaphore from "effect/Semaphore";
import * as SynchronizedRef from "effect/SynchronizedRef";

import * as ServerConfig from "../config.ts";
Expand Down Expand Up @@ -1569,7 +1569,7 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func
sessions: new Map(),
killFibers: new Map(),
});
const threadLocksRef = yield* SynchronizedRef.make(new Map<string, Semaphore.Semaphore>());
const threadLocks = yield* KeyedLock.make<string>();
const terminalEventListeners = new Set<(event: TerminalEvent) => Effect.Effect<void>>();
const workerScope = yield* Scope.make("sequential");
yield* Effect.addFinalizer(() => Scope.close(workerScope, Exit.void));
Expand Down Expand Up @@ -1598,29 +1598,10 @@ export const makeWithOptions = Effect.fn("TerminalManager.makeWithOptions")(func
f: (state: TerminalManagerState) => readonly [A, TerminalManagerState],
) => SynchronizedRef.modify(managerStateRef, f);

const getThreadSemaphore = (threadId: string) =>
SynchronizedRef.modifyEffect(threadLocksRef, (current) => {
const existing: Option.Option<Semaphore.Semaphore> = Option.fromNullishOr(
current.get(threadId),
);
return Option.match(existing, {
onNone: () =>
Semaphore.make(1).pipe(
Effect.map((semaphore) => {
const next = new Map(current);
next.set(threadId, semaphore);
return [semaphore, next] as const;
}),
),
onSome: (semaphore) => Effect.succeed([semaphore, current] as const),
});
});

const withThreadLock = <A, E, R>(
threadId: string,
effect: Effect.Effect<A, E, R>,
): Effect.Effect<A, E, R> =>
Effect.flatMap(getThreadSemaphore(threadId), (semaphore) => semaphore.withPermit(effect));
): Effect.Effect<A, E, R> => threadLocks.withLock(threadId, effect);

const clearKillFiber = Effect.fn("terminal.clearKillFiber")(function* (
process: PtyAdapter.PtyProcess | null,
Expand Down
14 changes: 4 additions & 10 deletions apps/server/src/vcs/VcsStatusBroadcaster.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@ import * as PubSub from "effect/PubSub";
import * as Ref from "effect/Ref";
import * as Schedule from "effect/Schedule";
import * as Scope from "effect/Scope";
import * as Semaphore from "effect/Semaphore";
import * as Stream from "effect/Stream";
import * as SynchronizedRef from "effect/SynchronizedRef";
import type {
Expand All @@ -22,6 +21,7 @@ import type {
VcsStatusStreamEvent,
} from "@t3tools/contracts";
import { mergeGitStatusParts } from "@t3tools/shared/git";
import * as KeyedLock from "@t3tools/shared/KeyedLock";
import { resolveProjectSettings } from "@t3tools/shared/projectSettings";

import * as BackgroundPolicy from "../background/BackgroundPolicy.ts";
Expand Down Expand Up @@ -234,15 +234,9 @@ export const make = Effect.gen(function* () {
// One permit per cwd for remote reads that write the cache. Without it a
// periodic poll that started before `gh pr create` can finish after the
// turn-end refresh and overwrite the fresh PR with its stale `pr: null`.
const remoteWriteLocks = new Map<string, Semaphore.Semaphore>();
const withRemoteWriteLock = <A, E, R>(cwd: string, effect: Effect.Effect<A, E, R>) => {
let lock = remoteWriteLocks.get(cwd);
if (lock === undefined) {
lock = Semaphore.makeUnsafe(1);
remoteWriteLocks.set(cwd, lock);
}
return lock.withPermits(1)(effect);
};
const remoteWriteLocks = yield* KeyedLock.make<string>();
const withRemoteWriteLock = <A, E, R>(cwd: string, effect: Effect.Effect<A, E, R>) =>
remoteWriteLocks.withLock(cwd, effect);
const pollersRef = yield* SynchronizedRef.make(new Map<string, ActiveRemotePoller>());

const getCachedStatus = Effect.fn("VcsStatusBroadcaster.getCachedStatus")(function* (
Expand Down
Loading
Loading