diff --git a/packages/mcp/bootstrap.ts b/packages/mcp/bootstrap.ts index 7deee12..e4586c3 100644 --- a/packages/mcp/bootstrap.ts +++ b/packages/mcp/bootstrap.ts @@ -571,7 +571,7 @@ export const deriveProject = ( deps: DeriveProjectDeps, ): Effect.Effect> => { 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, { @@ -604,7 +604,7 @@ const gitStdout = ( Command.stderr('pipe'), Command.string, Effect.map((out) => out.trim()), - Effect.catchAll(() => Effect.succeed('')), + Effect.orElseSucceed(() => ''), ) /** diff --git a/packages/mcp/channels-catch-up.test.ts b/packages/mcp/channels-catch-up.test.ts index 6778fb7..0e3c7f7 100644 --- a/packages/mcp/channels-catch-up.test.ts +++ b/packages/mcp/channels-catch-up.test.ts @@ -106,7 +106,7 @@ const buildHistorySpy = ( return byThread[`${channel}/${threadName}`] ?? [] }), recentThreads: () => Effect.succeed([]), - messagePermalink: () => Effect.succeed(Option.none()), + messagePermalink: () => Effect.succeedNone, }, } } diff --git a/packages/mcp/deployed-wiring.test.ts b/packages/mcp/deployed-wiring.test.ts index b3a9f5b..3ebf1a0 100644 --- a/packages/mcp/deployed-wiring.test.ts +++ b/packages/mcp/deployed-wiring.test.ts @@ -219,11 +219,11 @@ const bootDeployedSeat = async ( Deferred.unsafeDone(resumeOutcome, Effect.succeed(false)) const sessionIdDeferred = Deferred.unsafeMake(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, } diff --git a/packages/mcp/disconnect-exit.fixture.ts b/packages/mcp/disconnect-exit.fixture.ts index 4900a32..ef9eda0 100644 --- a/packages/mcp/disconnect-exit.fixture.ts +++ b/packages/mcp/disconnect-exit.fixture.ts @@ -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' @@ -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, } diff --git a/packages/mcp/queue-state-hooks.ts b/packages/mcp/queue-state-hooks.ts index c406c6c..1405f78 100644 --- a/packages/mcp/queue-state-hooks.ts +++ b/packages/mcp/queue-state-hooks.ts @@ -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)), ), }), diff --git a/packages/mcp/queue-state-store.ts b/packages/mcp/queue-state-store.ts index d607f17..a46216f 100644 --- a/packages/mcp/queue-state-store.ts +++ b/packages/mcp/queue-state-store.ts @@ -94,7 +94,7 @@ const readState = ( ): Effect.Effect, PlatformError | ParseResult.ParseError> => fs.readFileString(path).pipe( Effect.flatMap(decodeQueueStateFile), - Effect.map(Option.some), + Effect.asSome, Effect.catchIf(isNotFound, () => Effect.succeed(Option.none())), ) diff --git a/packages/mcp/server.integration.test.ts b/packages/mcp/server.integration.test.ts index 4067032..37fd587 100644 --- a/packages/mcp/server.integration.test.ts +++ b/packages/mcp/server.integration.test.ts @@ -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, } @@ -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, } @@ -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, } @@ -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[] = [] @@ -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 = { @@ -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(FiberId.none) const subscriptionStore: SubscriptionStore = { - read: () => Effect.succeed(Option.none()), + read: () => Effect.succeedNone, write: (intents) => Effect.flatMap(Deferred.await(sessionIdDeferred), (id) => Effect.sync(() => { @@ -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']) diff --git a/packages/mcp/server.test.ts b/packages/mcp/server.test.ts index 04301dc..3b87ade 100644 --- a/packages/mcp/server.test.ts +++ b/packages/mcp/server.test.ts @@ -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')), @@ -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 => Effect.succeed('offline'), } const adapter = completeAsSubstrate( @@ -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) }) diff --git a/packages/mcp/server.ts b/packages/mcp/server.ts index ce3c08b..f7c7274 100644 --- a/packages/mcp/server.ts +++ b/packages/mcp/server.ts @@ -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) @@ -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 => { - 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> => { + 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 @@ -593,12 +595,10 @@ export const makeProgram = ( Effect.succeedNone const rebuildNarrowSet: Effect.Effect = ( 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) => withSessionContext( restoreSubscriptions({ persisted: persistedTopicIntents, @@ -606,7 +606,7 @@ export const makeProgram = ( narrowSet, inbox: adapter.inbox, }), - { sessionId, project: parsed.project }, + { sessionId: Option.getOrUndefined(sessionId), project: parsed.project }, ), ), Effect.catchAll((err) => diff --git a/packages/mcp/subscription-restore.test.ts b/packages/mcp/subscription-restore.test.ts index 48ecfd8..c6c18c8 100644 --- a/packages/mcp/subscription-restore.test.ts +++ b/packages/mcp/subscription-restore.test.ts @@ -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' @@ -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 } @@ -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(() => { @@ -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 diff --git a/packages/mcp/subscription-store.ts b/packages/mcp/subscription-store.ts index 76b1b5f..8a43fdd 100644 --- a/packages/mcp/subscription-store.ts +++ b/packages/mcp/subscription-store.ts @@ -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>())), ) diff --git a/packages/mcp/tools-session.test.ts b/packages/mcp/tools-session.test.ts index c03e649..4006977 100644 --- a/packages/mcp/tools-session.test.ts +++ b/packages/mcp/tools-session.test.ts @@ -26,7 +26,7 @@ interface SessionRig { const buildSessionRig = ( options: { - readonly projectForCwd?: (cwd: string | undefined) => Effect.Effect + readonly projectForCwd?: (cwd: string | undefined) => Effect.Effect> readonly feedSessionId?: (sessionId: SessionId) => Effect.Effect } = {}, ): Effect.Effect => @@ -430,11 +430,9 @@ test('post with session_id + cwd mints cc-- 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({ @@ -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({ @@ -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', diff --git a/packages/mcp/tools.ts b/packages/mcp/tools.ts index 96e3a5b..894589a 100644 --- a/packages/mcp/tools.ts +++ b/packages/mcp/tools.ts @@ -194,10 +194,10 @@ export interface RegisterToolsDeps { * minted `cc--<8>` name reflects the *calling* session's * project rather than the plugin's own location. Wired at boot in * `server.ts` from `COMMY_PROJECT` (operator override) and - * the git probe; defaults to a constant `undefined` resolver when + * the git probe; defaults to a resolver that always answers none when * omitted (tests that don't care about per-session derivation). */ - readonly projectForCwd?: (cwd: string | undefined) => Effect.Effect + readonly projectForCwd?: (cwd: string | undefined) => Effect.Effect> /** * Restore (or seed) this session's narrow set on its first `subscribe`/ * `unsubscribe` — memoised once per session_id in `server.ts`, so it runs @@ -506,7 +506,7 @@ const UploadFileArgs = Schema.Struct({ path: Schema.String }) const buildToolDefs = (deps: RegisterToolsDeps, cache: InternalCache): ReadonlyArray => { const { adapter, identityCache, narrowSet } = deps - const projectForCwd = deps.projectForCwd ?? (() => Effect.succeed(undefined)) + const projectForCwd = deps.projectForCwd ?? (() => Effect.succeedNone) /** * The tools that accept a host-supplied `session_id` (see * {@link ToolDef.hostSuppliedArgs}). EIGHT tools carry it. That is the same @@ -548,7 +548,7 @@ const buildToolDefs = (deps: RegisterToolsDeps, cache: InternalCache): ReadonlyA } const projectForArgs = ( args: Readonly>, - ): Effect.Effect => projectForCwd(readCwd(args)) + ): Effect.Effect> => projectForCwd(readCwd(args)) // Every tool call that carries arguments supplies its calling session's // context, uniformly. This DECIDES NOTHING about identity: it feeds the // shared session-id deferred (comms-k7cv) and puts the naming inputs where @@ -569,7 +569,12 @@ const buildToolDefs = (deps: RegisterToolsDeps, cache: InternalCache): ReadonlyA projectForArgs(args).pipe( Effect.flatMap((project) => feedSession(sessionId).pipe( - Effect.zipRight(withSessionContext(effect, { sessionId, project })), + Effect.zipRight( + withSessionContext(effect, { + sessionId, + project: Option.getOrUndefined(project), + }), + ), ), ), ), @@ -937,7 +942,10 @@ const buildToolDefs = (deps: RegisterToolsDeps, cache: InternalCache): ReadonlyA // the snapshot persisted below captures the full live set and the // store's presence stays a true resume signal. if (sessionId !== undefined && deps.ensureSessionSubscriptions !== undefined) { - yield* deps.ensureSessionSubscriptions(sessionId, yield* projectForArgs(args)) + yield* deps.ensureSessionSubscriptions( + sessionId, + Option.getOrUndefined(yield* projectForArgs(args)), + ) } // Two sinks (see bootstrap.subscribeFromEnv): the consumer-side // narrow tells the event pump to tee matching events through; @@ -983,7 +991,10 @@ const buildToolDefs = (deps: RegisterToolsDeps, cache: InternalCache): ReadonlyA // Fed by `runFor` before this effect runs (see subscribe). const sessionId = readSessionId(args) if (sessionId !== undefined && deps.ensureSessionSubscriptions !== undefined) { - yield* deps.ensureSessionSubscriptions(sessionId, yield* projectForArgs(args)) + yield* deps.ensureSessionSubscriptions( + sessionId, + Option.getOrUndefined(yield* projectForArgs(args)), + ) } yield* Effect.sync(() => narrowSet.remove(intent)).pipe( Effect.andThen(adapter.inbox.unsubscribe(intentToTarget(intent))), diff --git a/packages/memory/adapter.test.ts b/packages/memory/adapter.test.ts index d32c8e9..73a451d 100644 --- a/packages/memory/adapter.test.ts +++ b/packages/memory/adapter.test.ts @@ -140,10 +140,9 @@ test('concurrent acquire on a fresh adapter binds exactly one name', async () => const adapter = await Effect.runPromise(memoryAdapter()) const names = ['agent-a', 'agent-b', 'agent-c', 'agent-d'].map((n) => decodeBotNameSync(n)) const exits = await Effect.runPromise( - Effect.all( - names.map((name) => Effect.exit(adapter.identity.acquire(name))), - { concurrency: 'unbounded' }, - ), + Effect.forEach(names, (name) => Effect.exit(adapter.identity.acquire(name)), { + concurrency: 'unbounded', + }), ) const successes = exits.filter(Exit.isSuccess) expect(successes).toHaveLength(1) diff --git a/packages/testing/contract.ts b/packages/testing/contract.ts index 1575d15..4323069 100644 --- a/packages/testing/contract.ts +++ b/packages/testing/contract.ts @@ -152,11 +152,7 @@ const takeUntil = ( queue: Queue.Queue, predicate: (event: InboundEvent) => boolean, ): Effect.Effect => - Queue.take(queue).pipe( - Effect.flatMap((event) => - predicate(event) ? Effect.succeed(event) : takeUntil(queue, predicate), - ), - ) + Queue.take(queue).pipe(Effect.filterOrElse(predicate, () => takeUntil(queue, predicate))) /** * Await the first event matching `predicate`, failing the test (with diff --git a/packages/zulip/adapter-events.test.ts b/packages/zulip/adapter-events.test.ts index e08e3c2..1429553 100644 --- a/packages/zulip/adapter-events.test.ts +++ b/packages/zulip/adapter-events.test.ts @@ -47,7 +47,7 @@ import { HttpClient } from '@effect/platform' import { Duration, Effect, - Option, + type Option, Queue, Redacted, type Scope, @@ -493,7 +493,7 @@ effectTest( Effect.gen(function* () { const stub = yield* makeStubHttpClient const adapter = yield* buildAdapterWithQueueConfig(stub, { - resumeQueue: () => Effect.succeed(Option.some({ queueId: 'resumed-q', lastEventId: 41 })), + resumeQueue: () => Effect.succeedSome({ queueId: 'resumed-q', lastEventId: 41 }), }) // The backlog buffered while the seat was dead: the reacted-to message // (which seeds the ref cache in-batch) followed by the reaction on it. @@ -569,7 +569,7 @@ effectTest( const stub = yield* makeStubHttpClient const outcomes: boolean[] = [] const adapter = yield* buildAdapterWithQueueConfig(stub, { - resumeQueue: () => Effect.succeed(Option.some({ queueId: 'resumed-q', lastEventId: 41 })), + resumeQueue: () => Effect.succeedSome({ queueId: 'resumed-q', lastEventId: 41 }), onResumeOutcome: (replayed) => Effect.sync(() => void outcomes.push(replayed)), }) yield* stub.respondSequence('GET', '/api/v1/events', [ @@ -594,7 +594,7 @@ effectTest( const stub = yield* makeStubHttpClient const outcomes: boolean[] = [] const adapter = yield* buildAdapterWithQueueConfig(stub, { - resumeQueue: () => Effect.succeed(Option.some({ queueId: 'dead-q', lastEventId: 41 })), + resumeQueue: () => Effect.succeedSome({ queueId: 'dead-q', lastEventId: 41 }), onResumeOutcome: (replayed) => Effect.sync(() => void outcomes.push(replayed)), }) // Fresh register for the re-registration after the dead resume-poll. diff --git a/packages/zulip/adapter.ts b/packages/zulip/adapter.ts index 4b80c41..4f595c5 100644 --- a/packages/zulip/adapter.ts +++ b/packages/zulip/adapter.ts @@ -1161,7 +1161,7 @@ export const zulipAdapter = ( Arr.findFirst(res.members, (u) => u.is_active && u.full_name === name).pipe( Option.match({ onNone: () => Effect.succeed(Option.none()), - onSome: (match) => toIdentity(match).pipe(Effect.map(Option.some)), + onSome: (match) => toIdentity(match).pipe(Effect.asSome), }), ), ), @@ -1339,7 +1339,10 @@ export const zulipAdapter = ( onSome: Effect.succeed, }), ), - Effect.flatMap((map) => (map.has(name) ? Effect.succeed(map) : refreshKnownStreams())), + Effect.filterOrElse( + (map) => map.has(name), + () => refreshKnownStreams(), + ), Effect.map((map) => Option.fromNullable(map.get(name))), ) @@ -1369,7 +1372,7 @@ export const zulipAdapter = ( Effect.flatMap((res) => Option.match(fromWireDescription(res.stream.description), { onNone: () => Effect.succeed(Option.none()), - onSome: (raw) => decodeChannelDescription(raw).pipe(Effect.map(Option.some)), + onSome: (raw) => decodeChannelDescription(raw).pipe(Effect.asSome), }), ), ) @@ -1470,7 +1473,7 @@ export const zulipAdapter = ( Arr.last(res.messages).pipe( Option.match({ onNone: () => Effect.succeed(Option.none()), - onSome: (m) => decodeMessageId(String(m.id)).pipe(Effect.map(Option.some)), + onSome: (m) => decodeMessageId(String(m.id)).pipe(Effect.asSome), }), ), ), @@ -2126,10 +2129,9 @@ export const zulipAdapter = ( { operator: 'topic', operand: topic }, ]) return readTopic(threadName).pipe( - Effect.flatMap((messages) => - messages.length > 0 - ? Effect.succeed(messages) - : readTopic(applyResolvedPrefix(threadName, true)), + Effect.filterOrElse( + (messages) => messages.length > 0, + () => readTopic(applyResolvedPrefix(threadName, true)), ), Effect.mapError((cause) => new HistoryError({ operation: 'readThread', cause })), ) diff --git a/packages/zulip/events.test.ts b/packages/zulip/events.test.ts index 2255573..7e5aa08 100644 --- a/packages/zulip/events.test.ts +++ b/packages/zulip/events.test.ts @@ -1958,7 +1958,7 @@ test( }, }), resolveDirectory: () => Effect.succeed(directoryFor(HERMES, MAINTAINER)), - currentRegistration: Effect.succeed(Option.some({ queueId: 'q-stale', lastEventId: 0 })), + currentRegistration: Effect.succeedSome({ queueId: 'q-stale', lastEventId: 0 }), boundIdentity: HERMES, messageRefCache: createMessageRefCache(), } diff --git a/packages/zulip/events.ts b/packages/zulip/events.ts index d1ab703..4a26ea5 100644 --- a/packages/zulip/events.ts +++ b/packages/zulip/events.ts @@ -586,8 +586,8 @@ export const fetchMessageRef = ( .pipe( Effect.flatMap((res) => { const message = res.messages[0] - if (message === undefined || message.id !== messageId) return Effect.succeed(Option.none()) - return decodeMessageRef(message, base).pipe(Effect.map(Option.some)) + if (message === undefined || message.id !== messageId) return Effect.succeedNone + return decodeMessageRef(message, base).pipe(Effect.asSome) }), ) diff --git a/packages/zulip/http.ts b/packages/zulip/http.ts index b988695..3bc5f6b 100644 --- a/packages/zulip/http.ts +++ b/packages/zulip/http.ts @@ -372,8 +372,10 @@ export const makeZulipHttp = ( Effect.flatMap((text) => classifyEnvelope(text, response.status, url)), ), ), - Effect.catchTag('RequestError', (cause) => Effect.fail(transportError(url, cause))), - Effect.catchTag('ResponseError', (cause) => Effect.fail(transportError(url, cause))), + Effect.catchTags({ + RequestError: (cause) => Effect.fail(transportError(url, cause)), + ResponseError: (cause) => Effect.fail(transportError(url, cause)), + }), ) const sendWithRetry = ( @@ -465,10 +467,10 @@ export const makeZulipHttp = ( Effect.mapError((cause) => transportError(url, cause)), ) }), - Effect.catchTag('RequestError', (cause) => Effect.fail(transportError(url, cause))), - Effect.catchTag('ResponseError', (cause) => - Effect.fail(transportError(url, cause)), - ), + Effect.catchTags({ + RequestError: (cause) => Effect.fail(transportError(url, cause)), + ResponseError: (cause) => Effect.fail(transportError(url, cause)), + }), ) })() : Effect.die(