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
4 changes: 2 additions & 2 deletions packages/mcp/bootstrap.ts
Original file line number Diff line number Diff line change
Expand Up @@ -571,7 +571,7 @@ export const deriveProject = (
deps: DeriveProjectDeps,
): Effect.Effect<Option.Option<ProjectSlug>> => {
if (deps.project !== undefined) {
return Effect.succeed(Option.some(deps.project))
return Effect.succeedSome(deps.project)
}
return Effect.map(deps.readGitContext(deps.cwd), (context) =>
matchGitContext(context, {
Expand Down Expand Up @@ -604,7 +604,7 @@ const gitStdout = (
Command.stderr('pipe'),
Command.string,
Effect.map((out) => out.trim()),
Effect.catchAll(() => Effect.succeed('')),
Effect.orElseSucceed(() => ''),
)

/**
Expand Down
2 changes: 1 addition & 1 deletion packages/mcp/channels-catch-up.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ const buildHistorySpy = (
return byThread[`${channel}/${threadName}`] ?? []
}),
recentThreads: () => Effect.succeed([]),
messagePermalink: () => Effect.succeed(Option.none()),
messagePermalink: () => Effect.succeedNone,
},
}
}
Expand Down
4 changes: 2 additions & 2 deletions packages/mcp/deployed-wiring.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -219,11 +219,11 @@ const bootDeployedSeat = async (
Deferred.unsafeDone(resumeOutcome, Effect.succeed(false))
const sessionIdDeferred = Deferred.unsafeMake<SessionIdValue>(FiberId.none)
const inMemoryCursorStore = {
read: () => Effect.succeed(Option.none()),
read: () => Effect.succeedNone,
write: () => Effect.void,
}
const inMemorySubscriptionStore = {
read: () => Effect.succeed(Option.none()),
read: () => Effect.succeedNone,
write: () => Effect.void,
}

Expand Down
6 changes: 3 additions & 3 deletions packages/mcp/disconnect-exit.fixture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ import { memoryAdapter } from '@commy/memory/adapter'
import { FetchHttpClient } from '@effect/platform'
import { NodeContext, NodeRuntime } from '@effect/platform-node'
import { StdioServerTransport } from '@modelcontextprotocol/sdk/server/stdio.js'
import { ConfigProvider, Effect, Layer, Option } from 'effect'
import { ConfigProvider, Effect, Layer } from 'effect'
import { substrateAdapterLayer } from './bootstrap.ts'
import { CursorStoreTag } from './cursor-store.ts'
import { completeAsSubstrate } from './memory-substrate.ts'
Expand All @@ -33,12 +33,12 @@ import { SubscriptionStoreTag } from './subscription-store.ts'
import { testBootStoresLayer } from './test-platform.ts'

const inMemoryCursorStore = {
read: () => Effect.succeed(Option.none()),
read: () => Effect.succeedNone,
write: () => Effect.void,
}

const inMemorySubscriptionStore = {
read: () => Effect.succeed(Option.none()),
read: () => Effect.succeedNone,
write: () => Effect.void,
}

Expand Down
2 changes: 1 addition & 1 deletion packages/mcp/queue-state-hooks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,7 @@ export const buildQueueStateHooks = (deps: {
onNone: () =>
store.read(id).pipe(
Effect.map(Option.isNone),
Effect.catchAll(() => Effect.succeed(true)),
Effect.orElseSucceed(() => true),
Effect.flatMap((isFresh) => (isFresh ? store.write(id, queue) : Effect.void)),
),
}),
Expand Down
2 changes: 1 addition & 1 deletion packages/mcp/queue-state-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ const readState = (
): Effect.Effect<Option.Option<EventQueueCursor>, PlatformError | ParseResult.ParseError> =>
fs.readFileString(path).pipe(
Effect.flatMap(decodeQueueStateFile),
Effect.map(Option.some),
Effect.asSome,
Effect.catchIf(isNotFound, () => Effect.succeed(Option.none<EventQueueCursor>())),
)

Expand Down
14 changes: 7 additions & 7 deletions packages/mcp/server.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1367,7 +1367,7 @@ test('persistent boot with no prior cursor: no replay events fired, cursor initi
test('persistent boot with a prior cursor: replay fires, mention-received notifications dispatched ahead of pump', async () => {
const PRIOR_CURSOR_TS = 1000
const cursorStore: CursorStore = {
read: () => Effect.succeed(Option.some(decodeTimestampSync(PRIOR_CURSOR_TS))),
read: () => Effect.succeedSome(decodeTimestampSync(PRIOR_CURSOR_TS)),
write: () => Effect.void,
}

Expand Down Expand Up @@ -1529,7 +1529,7 @@ test('ephemeral lazy acquire with no prior cursor: no replay, cursor initialised
test('ephemeral lazy acquire with a prior cursor: replay fires, mention dispatched ahead of tool result', async () => {
const PRIOR_CURSOR_TS = 1000
const cursorStore: CursorStore = {
read: () => Effect.succeed(Option.some(decodeTimestampSync(PRIOR_CURSOR_TS))),
read: () => Effect.succeedSome(decodeTimestampSync(PRIOR_CURSOR_TS)),
write: () => Effect.void,
}

Expand Down Expand Up @@ -1615,7 +1615,7 @@ test('ephemeral lazy acquire with a prior cursor: replay fires, mention dispatch

test('ephemeral catch-up failure is non-fatal: tool call succeeds, failure is logged', async () => {
const cursorStore: CursorStore = {
read: () => Effect.succeed(Option.some(decodeTimestampSync(1000))),
read: () => Effect.succeedSome(decodeTimestampSync(1000)),
write: () => Effect.void,
}

Expand Down Expand Up @@ -1664,7 +1664,7 @@ test('ephemeral queue-ALIVE resume: history catch-up does NOT run, no duplicate
// either would double-deliver every message the pump already carries.
const PRIOR_CURSOR_TS = 1000
const cursorStore: CursorStore = {
read: () => Effect.succeed(Option.some(decodeTimestampSync(PRIOR_CURSOR_TS))),
read: () => Effect.succeedSome(decodeTimestampSync(PRIOR_CURSOR_TS)),
write: () => Effect.void,
}
const replayCalls: number[] = []
Expand Down Expand Up @@ -1720,7 +1720,7 @@ test('ephemeral queue-DEAD resume: mentions + channels catch-up run and backfill
// messages the dead queue could not carry.
const PRIOR_CURSOR_TS = 1000
const cursorStore: CursorStore = {
read: () => Effect.succeed(Option.some(decodeTimestampSync(PRIOR_CURSOR_TS))),
read: () => Effect.succeedSome(decodeTimestampSync(PRIOR_CURSOR_TS)),
write: () => Effect.void,
}
const mentionedIdentity = {
Expand Down Expand Up @@ -1997,7 +1997,7 @@ test('ephemeral subscribe persists the live narrow set (defaults + new sub) unde
// keyed under that id with no id ever passed to write().
const sessionIdDeferred = Deferred.unsafeMake<SessionIdValue>(FiberId.none)
const subscriptionStore: SubscriptionStore = {
read: () => Effect.succeed(Option.none()),
read: () => Effect.succeedNone,
write: (intents) =>
Effect.flatMap(Deferred.await(sessionIdDeferred), (id) =>
Effect.sync(() => {
Expand Down Expand Up @@ -2035,7 +2035,7 @@ test('ephemeral resume recovers its channels from the realm and does NOT re-appl
// is what marks this a resume rather than a first launch. Resume must honour
// both: recover `home`, and never re-add a dropped default.
const subscriptionStore: SubscriptionStore = {
read: () => Effect.succeed(Option.some([])),
read: () => Effect.succeedSome([]),
write: () => Effect.void,
}
const cap = captureSubscribes(['home'])
Expand Down
32 changes: 15 additions & 17 deletions packages/mcp/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,7 @@ const buildFakeAdapter = (
return Effect.succeed(acquiredIdentity)
}),
release: () => Effect.void,
resolve: () => Effect.succeed(Option.none()),
resolve: () => Effect.succeedNone,
}
const publisher: MessagePublisher = {
post: () => Effect.die(new Error('unused fake')),
Expand Down Expand Up @@ -247,13 +247,13 @@ const buildFakeAdapter = (
readChannel: (_channel: ChannelName, _range: Range) => Effect.succeed([]),
readThread: (_channel: ChannelName, _threadName, _range?: Range) => Effect.succeed([]),
recentThreads: () => Effect.succeed([]),
messagePermalink: () => Effect.succeed(Option.none()),
messagePermalink: () => Effect.succeedNone,
}
const directory: Directory = {
listAgents: () => Effect.succeed([]),
listHumans: () => Effect.succeed([]),
listChannels: () => Effect.succeed([]),
channelDescription: () => Effect.succeed(Option.none()),
channelDescription: () => Effect.succeedNone,
presence: (_id: Identity): Effect.Effect<Presence> => Effect.succeed('offline'),
}
const adapter = completeAsSubstrate(
Expand Down Expand Up @@ -532,20 +532,18 @@ test('boot completes for an ephemeral seat with COMMY_SUBSCRIBE and no resume ve
{ ...adapter, inbox: { ...adapter.inbox, events: neverProduces } },
{ close: async () => {} },
)
const exit = await Effect.runPromise(
Effect.exit(
Effect.promise(() =>
runProgram(
{ ...lazyEnv, COMMY_SUBSCRIBE: 'home' },
substrate,
{
loggerLayer: captureLogger([]),
readGitContext: () => Effect.succeed(NotInRepo()),
},
binderRef,
),
).pipe(Effect.timeoutFail({ duration: '5 seconds', onTimeout: () => 'boot hung' as const })),
),
const exit = await Effect.runPromiseExit(
Effect.promise(() =>
runProgram(
{ ...lazyEnv, COMMY_SUBSCRIBE: 'home' },
substrate,
{
loggerLayer: captureLogger([]),
readGitContext: () => Effect.succeed(NotInRepo()),
},
binderRef,
),
).pipe(Effect.timeoutFail({ duration: '5 seconds', onTimeout: () => 'boot hung' as const })),
)
expect(Exit.isSuccess(exit)).toBe(true)
})
Expand Down
22 changes: 11 additions & 11 deletions packages/mcp/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -456,7 +456,7 @@ export const makeProgram = (
Effect.map(Option.isNone),
// An unreadable store is treated as fresh, matching
// `resumeQueue`'s own best-effort degrade.
Effect.catchAll(() => Effect.succeed(true)),
Effect.orElseSucceed(() => true),
Effect.flatMap((nothingPersisted) =>
nothingPersisted
? Deferred.succeed(resumeOutcome, false).pipe(Effect.asVoid)
Expand Down Expand Up @@ -524,10 +524,12 @@ export const makeProgram = (
// Per-call project resolver. Operator override (COMMY_PROJECT)
// is authoritative; otherwise derive from the calling session's cwd
// at call time — process cwd is irrelevant.
const projectForCwd = (cwd: string | undefined): Effect.Effect<ProjectSlug | undefined> => {
if (parsed.project !== undefined) return Effect.succeed(parsed.project)
if (cwd === undefined) return Effect.succeed(undefined)
return Effect.map(deriveProject({ cwd, readGitContext }), Option.getOrUndefined)
const projectForCwd = (
cwd: string | undefined,
): Effect.Effect<Option.Option<ProjectSlug>> => {
if (parsed.project !== undefined) return Effect.succeedSome(parsed.project)
if (cwd === undefined) return Effect.succeedNone
return deriveProject({ cwd, readGitContext })
}

// Sample the realm-wide editing switch once, before the tool list is
Expand Down Expand Up @@ -593,20 +595,18 @@ export const makeProgram = (
Effect.succeedNone
const rebuildNarrowSet: Effect.Effect<void> = (
parsed.botName === undefined
? Deferred.await(sessionIdDeferred).pipe(
Effect.map((sessionId): SessionId | undefined => sessionId),
)
: Effect.succeed(undefined)
? Deferred.await(sessionIdDeferred).pipe(Effect.asSome)
: Effect.succeedNone
).pipe(
Effect.flatMap((sessionId) =>
Effect.flatMap((sessionId: Option.Option<SessionId>) =>
withSessionContext(
restoreSubscriptions({
persisted: persistedTopicIntents,
isBound: () => identityCache.boundIdentityIds().size > 0,
narrowSet,
inbox: adapter.inbox,
}),
{ sessionId, project: parsed.project },
{ sessionId: Option.getOrUndefined(sessionId), project: parsed.project },
),
),
Effect.catchAll((err) =>
Expand Down
8 changes: 4 additions & 4 deletions packages/mcp/subscription-restore.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { describe, expect, test } from 'bun:test'
import { decodeChannelNameSync } from '@commy/core/ports'
import { Effect, Option } from 'effect'
import { Effect } from 'effect'
import type { ProjectSlug } from './bootstrap.ts'
import { createNarrowSet } from './narrow-set.ts'
import type { SubscribeIntent } from './subscribe-parser.ts'
Expand Down Expand Up @@ -35,7 +35,7 @@ describe('seedDefaultsIfFresh', () => {
await Effect.runPromise(
seedDefaultsIfFresh(
{
subscriptionStore: stubStore(() => Effect.succeed(Option.none())),
subscriptionStore: stubStore(() => Effect.succeedNone),
registerDefaults: (project) =>
Effect.sync(() => {
defaultsCall = { project }
Expand All @@ -53,7 +53,7 @@ describe('seedDefaultsIfFresh', () => {
seedDefaultsIfFresh(
{
subscriptionStore: stubStore(() =>
Effect.succeed(Option.some([newTopics('general'), channel('commy')])),
Effect.succeedSome([newTopics('general'), channel('commy')]),
),
registerDefaults: () =>
Effect.sync(() => {
Expand All @@ -71,7 +71,7 @@ describe('seedDefaultsIfFresh', () => {
await Effect.runPromise(
seedDefaultsIfFresh(
{
subscriptionStore: stubStore(() => Effect.succeed(Option.some([]))),
subscriptionStore: stubStore(() => Effect.succeedSome([])),
registerDefaults: () =>
Effect.sync(() => {
defaultsCalled = true
Expand Down
2 changes: 1 addition & 1 deletion packages/mcp/subscription-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,7 @@ const readSubscriptions = (
> =>
fs.readFileString(path).pipe(
Effect.flatMap((raw) => decodeSubscriptionsFile(id, raw)),
Effect.map((intents) => Option.some(intents)),
Effect.asSome,
Effect.catchIf(isNotFound, () => Effect.succeed(Option.none<ReadonlyArray<SubscribeIntent>>())),
)

Expand Down
17 changes: 8 additions & 9 deletions packages/mcp/tools-session.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ interface SessionRig {

const buildSessionRig = (
options: {
readonly projectForCwd?: (cwd: string | undefined) => Effect.Effect<ProjectSlug | undefined>
readonly projectForCwd?: (cwd: string | undefined) => Effect.Effect<Option.Option<ProjectSlug>>
readonly feedSessionId?: (sessionId: SessionId) => Effect.Effect<void>
} = {},
): Effect.Effect<SessionRig, never, Scope.Scope> =>
Expand Down Expand Up @@ -430,11 +430,9 @@ test('post with session_id + cwd mints cc-<project>-<sid-prefix> derived from cw
Effect.gen(function* () {
const rig = yield* buildSessionRig({
projectForCwd: (cwd) =>
Effect.succeed(
cwd === '/home/x/myproject'
? Option.getOrUndefined(sanitiseProjectSlug('myproject'))
: undefined,
),
cwd === '/home/x/myproject'
? Effect.succeed(sanitiseProjectSlug('myproject'))
: Effect.succeedNone,
})
const result = yield* Effect.promise(() =>
rig.client.callTool({
Expand Down Expand Up @@ -466,7 +464,8 @@ test('two sessions in different cwds mint two different project prefixes', () =>
'/home/x/myproject-b': slug('myproject-b'),
}
const rig = yield* buildSessionRig({
projectForCwd: (cwd) => Effect.succeed(cwd === undefined ? undefined : cwdToSlug[cwd]),
projectForCwd: (cwd) =>
Effect.succeed(Option.fromNullable(cwd === undefined ? undefined : cwdToSlug[cwd])),
})
yield* Effect.promise(() =>
rig.client.callTool({
Expand Down Expand Up @@ -503,9 +502,9 @@ test('post with cwd from a non-project directory falls back to bare cc-<8>', ()
Effect.runPromise(
Effect.scoped(
Effect.gen(function* () {
// projectForCwd returns undefined when cwd is not in a known repo.
// projectForCwd returns none when cwd is not in a known repo.
// The minted name must NOT inherit the plugin's own location.
const rig = yield* buildSessionRig({ projectForCwd: () => Effect.succeed(undefined) })
const rig = yield* buildSessionRig({ projectForCwd: () => Effect.succeedNone })
yield* Effect.promise(() =>
rig.client.callTool({
name: 'post',
Expand Down
Loading