test(server): migrate session messages to Effect runner (#27234)
This commit is contained in:
@@ -1,191 +1,179 @@
|
|||||||
import { afterEach, describe, expect, test } from "bun:test"
|
import { afterEach, describe, expect } from "bun:test"
|
||||||
import { Effect } from "effect"
|
import { Effect } from "effect"
|
||||||
import { WithInstance } from "../../src/project/with-instance"
|
|
||||||
import { Server } from "../../src/server/server"
|
import { Server } from "../../src/server/server"
|
||||||
import { Session as SessionNs } from "@/session/session"
|
import { Session as SessionNs } from "@/session/session"
|
||||||
import { MessageV2 } from "../../src/session/message-v2"
|
import { MessageV2 } from "../../src/session/message-v2"
|
||||||
|
import { ModelID, ProviderID } from "../../src/provider/schema"
|
||||||
import { MessageID, PartID, type SessionID } from "../../src/session/schema"
|
import { MessageID, PartID, type SessionID } from "../../src/session/schema"
|
||||||
import * as Log from "@opencode-ai/core/util/log"
|
import * as Log from "@opencode-ai/core/util/log"
|
||||||
import { disposeAllInstances, tmpdir } from "../fixture/fixture"
|
import { disposeAllInstances, TestInstance } from "../fixture/fixture"
|
||||||
|
import { testEffect } from "../lib/effect"
|
||||||
|
|
||||||
void Log.init({ print: false })
|
void Log.init({ print: false })
|
||||||
|
|
||||||
function run<A, E>(fx: Effect.Effect<A, E, SessionNs.Service>) {
|
const it = testEffect(SessionNs.defaultLayer)
|
||||||
return Effect.runPromise(fx.pipe(Effect.provide(SessionNs.defaultLayer)))
|
|
||||||
}
|
|
||||||
|
|
||||||
const svc = {
|
const model = {
|
||||||
...SessionNs,
|
providerID: ProviderID.make("test"),
|
||||||
create(input?: SessionNs.CreateInput) {
|
modelID: ModelID.make("test"),
|
||||||
return run(SessionNs.Service.use((svc) => svc.create(input)))
|
|
||||||
},
|
|
||||||
remove(id: SessionID) {
|
|
||||||
return run(SessionNs.Service.use((svc) => svc.remove(id)))
|
|
||||||
},
|
|
||||||
updateMessage<T extends MessageV2.Info>(msg: T) {
|
|
||||||
return run(SessionNs.Service.use((svc) => svc.updateMessage(msg)))
|
|
||||||
},
|
|
||||||
updatePart<T extends MessageV2.Part>(part: T) {
|
|
||||||
return run(SessionNs.Service.use((svc) => svc.updatePart(part)))
|
|
||||||
},
|
|
||||||
}
|
}
|
||||||
|
|
||||||
afterEach(async () => {
|
afterEach(async () => {
|
||||||
await disposeAllInstances()
|
await disposeAllInstances()
|
||||||
})
|
})
|
||||||
|
|
||||||
async function withoutWatcher<T>(fn: () => Promise<T>) {
|
const withoutWatcher = <A, E, R>(effect: Effect.Effect<A, E, R>) => {
|
||||||
if (process.platform !== "win32") return fn()
|
if (process.platform !== "win32") return effect
|
||||||
const prev = process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER
|
return Effect.acquireUseRelease(
|
||||||
process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER = "true"
|
Effect.sync(() => {
|
||||||
try {
|
const previous = process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER
|
||||||
return await fn()
|
process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER = "true"
|
||||||
} finally {
|
return previous
|
||||||
if (prev === undefined) delete process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER
|
}),
|
||||||
else process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER = prev
|
() => effect,
|
||||||
}
|
(previous) =>
|
||||||
|
Effect.sync(() => {
|
||||||
|
if (previous === undefined) delete process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER
|
||||||
|
else process.env.OPENCODE_EXPERIMENTAL_DISABLE_FILEWATCHER = previous
|
||||||
|
}),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
async function fill(sessionID: SessionID, count: number, time = (i: number) => Date.now() + i) {
|
const sessionScoped = Effect.acquireRelease(
|
||||||
const ids = [] as MessageID[]
|
SessionNs.Service.use((svc) => svc.create({})),
|
||||||
for (let i = 0; i < count; i++) {
|
(session) => SessionNs.Service.use((svc) => svc.remove(session.id)).pipe(Effect.ignore),
|
||||||
const id = MessageID.ascending()
|
)
|
||||||
ids.push(id)
|
|
||||||
await svc.updateMessage({
|
const fill = Effect.fn("SessionMessagesTest.fill")(function* (
|
||||||
id,
|
sessionID: SessionID,
|
||||||
sessionID,
|
count: number,
|
||||||
role: "user",
|
time = (i: number) => Date.now() + i,
|
||||||
time: { created: time(i) },
|
) {
|
||||||
agent: "test",
|
const session = yield* SessionNs.Service
|
||||||
model: { providerID: "test", modelID: "test" },
|
return yield* Effect.forEach(
|
||||||
tools: {},
|
Array.from({ length: count }, (_, i) => i),
|
||||||
mode: "",
|
(i) =>
|
||||||
} as unknown as MessageV2.Info)
|
Effect.gen(function* () {
|
||||||
await svc.updatePart({
|
const id = MessageID.ascending()
|
||||||
id: PartID.ascending(),
|
yield* session.updateMessage({
|
||||||
sessionID,
|
id,
|
||||||
messageID: id,
|
sessionID,
|
||||||
type: "text",
|
role: "user",
|
||||||
text: `m${i}`,
|
time: { created: time(i) },
|
||||||
})
|
agent: "test",
|
||||||
}
|
model,
|
||||||
return ids
|
tools: {},
|
||||||
|
} satisfies MessageV2.User)
|
||||||
|
yield* session.updatePart({
|
||||||
|
id: PartID.ascending(),
|
||||||
|
sessionID,
|
||||||
|
messageID: id,
|
||||||
|
type: "text",
|
||||||
|
text: `m${i}`,
|
||||||
|
} satisfies MessageV2.TextPart)
|
||||||
|
return id
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
function request(path: string) {
|
||||||
|
return Effect.promise(() => Promise.resolve(Server.Default().app.request(path)))
|
||||||
|
}
|
||||||
|
|
||||||
|
function json<T>(response: Response) {
|
||||||
|
return Effect.promise(() => response.json() as Promise<T>)
|
||||||
}
|
}
|
||||||
|
|
||||||
describe("session messages endpoint", () => {
|
describe("session messages endpoint", () => {
|
||||||
test("returns cursor headers for older pages", async () => {
|
it.instance(
|
||||||
await using tmp = await tmpdir({ git: true })
|
"returns cursor headers for older pages",
|
||||||
await withoutWatcher(() =>
|
withoutWatcher(
|
||||||
WithInstance.provide({
|
Effect.gen(function* () {
|
||||||
directory: tmp.path,
|
const session = yield* sessionScoped
|
||||||
fn: async () => {
|
const ids = yield* fill(session.id, 5)
|
||||||
const session = await svc.create({})
|
|
||||||
const ids = await fill(session.id, 5)
|
|
||||||
const app = Server.Default().app
|
|
||||||
|
|
||||||
const a = await app.request(`/session/${session.id}/message?limit=2`)
|
const a = yield* request(`/session/${session.id}/message?limit=2`)
|
||||||
expect(a.status).toBe(200)
|
expect(a.status).toBe(200)
|
||||||
const aBody = (await a.json()) as MessageV2.WithParts[]
|
const aBody = yield* json<MessageV2.WithParts[]>(a)
|
||||||
expect(aBody.map((item) => item.info.id)).toEqual(ids.slice(-2))
|
expect(aBody.map((item) => item.info.id)).toEqual(ids.slice(-2))
|
||||||
const cursor = a.headers.get("x-next-cursor")
|
const cursor = a.headers.get("x-next-cursor")
|
||||||
expect(cursor).toBeTruthy()
|
expect(cursor).toBeTruthy()
|
||||||
expect(a.headers.get("link")).toContain('rel="next"')
|
expect(a.headers.get("link")).toContain('rel="next"')
|
||||||
|
|
||||||
const b = await app.request(`/session/${session.id}/message?limit=2&before=${encodeURIComponent(cursor!)}`)
|
const b = yield* request(`/session/${session.id}/message?limit=2&before=${encodeURIComponent(cursor!)}`)
|
||||||
expect(b.status).toBe(200)
|
expect(b.status).toBe(200)
|
||||||
const bBody = (await b.json()) as MessageV2.WithParts[]
|
const bBody = yield* json<MessageV2.WithParts[]>(b)
|
||||||
expect(bBody.map((item) => item.info.id)).toEqual(ids.slice(-4, -2))
|
expect(bBody.map((item) => item.info.id)).toEqual(ids.slice(-4, -2))
|
||||||
|
|
||||||
await svc.remove(session.id)
|
|
||||||
},
|
|
||||||
}),
|
}),
|
||||||
)
|
),
|
||||||
})
|
{ git: true },
|
||||||
|
)
|
||||||
|
|
||||||
test("keeps full-history responses when limit is omitted", async () => {
|
it.instance(
|
||||||
await using tmp = await tmpdir({ git: true })
|
"keeps full-history responses when limit is omitted",
|
||||||
await withoutWatcher(() =>
|
withoutWatcher(
|
||||||
WithInstance.provide({
|
Effect.gen(function* () {
|
||||||
directory: tmp.path,
|
const session = yield* sessionScoped
|
||||||
fn: async () => {
|
const ids = yield* fill(session.id, 3)
|
||||||
const session = await svc.create({})
|
|
||||||
const ids = await fill(session.id, 3)
|
|
||||||
const app = Server.Default().app
|
|
||||||
|
|
||||||
const res = await app.request(`/session/${session.id}/message`)
|
const res = yield* request(`/session/${session.id}/message`)
|
||||||
expect(res.status).toBe(200)
|
expect(res.status).toBe(200)
|
||||||
const body = (await res.json()) as MessageV2.WithParts[]
|
const body = yield* json<MessageV2.WithParts[]>(res)
|
||||||
expect(body.map((item) => item.info.id)).toEqual(ids)
|
expect(body.map((item) => item.info.id)).toEqual(ids)
|
||||||
|
|
||||||
await svc.remove(session.id)
|
|
||||||
},
|
|
||||||
}),
|
}),
|
||||||
)
|
),
|
||||||
})
|
{ git: true },
|
||||||
|
)
|
||||||
|
|
||||||
test("rejects invalid cursors and missing sessions", async () => {
|
it.instance(
|
||||||
await using tmp = await tmpdir({ git: true })
|
"rejects invalid cursors and missing sessions",
|
||||||
await withoutWatcher(() =>
|
withoutWatcher(
|
||||||
WithInstance.provide({
|
Effect.gen(function* () {
|
||||||
directory: tmp.path,
|
const session = yield* sessionScoped
|
||||||
fn: async () => {
|
|
||||||
const session = await svc.create({})
|
|
||||||
const app = Server.Default().app
|
|
||||||
|
|
||||||
const bad = await app.request(`/session/${session.id}/message?limit=2&before=bad`)
|
const bad = yield* request(`/session/${session.id}/message?limit=2&before=bad`)
|
||||||
expect(bad.status).toBe(400)
|
expect(bad.status).toBe(400)
|
||||||
|
|
||||||
const miss = await app.request(`/session/ses_missing/message?limit=2`)
|
const miss = yield* request(`/session/ses_missing/message?limit=2`)
|
||||||
expect(miss.status).toBe(404)
|
expect(miss.status).toBe(404)
|
||||||
|
|
||||||
await svc.remove(session.id)
|
|
||||||
},
|
|
||||||
}),
|
}),
|
||||||
)
|
),
|
||||||
})
|
{ git: true },
|
||||||
|
)
|
||||||
|
|
||||||
test("does not truncate large legacy limit requests", async () => {
|
it.instance(
|
||||||
await using tmp = await tmpdir({ git: true })
|
"does not truncate large legacy limit requests",
|
||||||
await withoutWatcher(() =>
|
withoutWatcher(
|
||||||
WithInstance.provide({
|
Effect.gen(function* () {
|
||||||
directory: tmp.path,
|
const session = yield* sessionScoped
|
||||||
fn: async () => {
|
yield* fill(session.id, 520)
|
||||||
const session = await svc.create({})
|
|
||||||
await fill(session.id, 520)
|
|
||||||
const app = Server.Default().app
|
|
||||||
|
|
||||||
const res = await app.request(`/session/${session.id}/message?limit=510`)
|
const res = yield* request(`/session/${session.id}/message?limit=510`)
|
||||||
expect(res.status).toBe(200)
|
expect(res.status).toBe(200)
|
||||||
const body = (await res.json()) as MessageV2.WithParts[]
|
const body = yield* json<MessageV2.WithParts[]>(res)
|
||||||
expect(body).toHaveLength(510)
|
expect(body).toHaveLength(510)
|
||||||
|
|
||||||
await svc.remove(session.id)
|
|
||||||
},
|
|
||||||
}),
|
}),
|
||||||
)
|
),
|
||||||
})
|
{ git: true },
|
||||||
|
)
|
||||||
|
|
||||||
test("accepts directory query used by workspace routing", async () => {
|
it.instance(
|
||||||
await using tmp = await tmpdir({ git: true })
|
"accepts directory query used by workspace routing",
|
||||||
await withoutWatcher(() =>
|
withoutWatcher(
|
||||||
WithInstance.provide({
|
Effect.gen(function* () {
|
||||||
directory: tmp.path,
|
const tmp = yield* TestInstance
|
||||||
fn: async () => {
|
const session = yield* sessionScoped
|
||||||
const session = await svc.create({})
|
yield* fill(session.id, 1)
|
||||||
await fill(session.id, 1)
|
|
||||||
const app = Server.Default().app
|
|
||||||
|
|
||||||
const res = await app.request(
|
const res = yield* request(
|
||||||
`/session/${session.id}/message?limit=80&directory=${encodeURIComponent(tmp.path)}`,
|
`/session/${session.id}/message?limit=80&directory=${encodeURIComponent(tmp.directory)}`,
|
||||||
)
|
)
|
||||||
expect(res.status).toBe(200)
|
expect(res.status).toBe(200)
|
||||||
const body = await res.json()
|
const body = yield* json<unknown[]>(res)
|
||||||
expect(Array.isArray(body)).toBe(true)
|
expect(Array.isArray(body)).toBe(true)
|
||||||
expect(body).toHaveLength(1)
|
expect(body).toHaveLength(1)
|
||||||
|
|
||||||
await svc.remove(session.id)
|
|
||||||
},
|
|
||||||
}),
|
}),
|
||||||
)
|
),
|
||||||
})
|
{ git: true },
|
||||||
|
)
|
||||||
})
|
})
|
||||||
|
|||||||
Reference in New Issue
Block a user