fix: synchronize remote data refreshes
This commit is contained in:
parent
cc686ab8f6
commit
33c0cc2bb6
4 changed files with 36 additions and 20 deletions
|
|
@ -194,7 +194,7 @@ export const OpencodePlugin = define<HttpClient.HttpClient | EventV2.Service | S
|
||||||
Stream.runForEach(refresh),
|
Stream.runForEach(refresh),
|
||||||
Effect.forkScoped({ startImmediately: true }),
|
Effect.forkScoped({ startImmediately: true }),
|
||||||
)
|
)
|
||||||
yield* refresh().pipe(Effect.forkScoped)
|
yield* refresh()
|
||||||
}),
|
}),
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -163,14 +163,11 @@ describe("OpencodePlugin", () => {
|
||||||
Effect.acquireUseRelease(
|
Effect.acquireUseRelease(
|
||||||
Effect.sync(() => {
|
Effect.sync(() => {
|
||||||
const authorization: Array<string | null> = []
|
const authorization: Array<string | null> = []
|
||||||
const gate = Promise.withResolvers<void>()
|
|
||||||
return {
|
return {
|
||||||
authorization,
|
authorization,
|
||||||
release: gate.resolve,
|
|
||||||
server: Bun.serve({
|
server: Bun.serve({
|
||||||
port: 0,
|
port: 0,
|
||||||
fetch: async (request) => {
|
fetch: (request) => {
|
||||||
await gate.promise
|
|
||||||
authorization.push(request.headers.get("authorization"))
|
authorization.push(request.headers.get("authorization"))
|
||||||
const origin = new URL(request.url).origin
|
const origin = new URL(request.url).origin
|
||||||
return Response.json({
|
return Response.json({
|
||||||
|
|
@ -209,7 +206,7 @@ describe("OpencodePlugin", () => {
|
||||||
}),
|
}),
|
||||||
}
|
}
|
||||||
}),
|
}),
|
||||||
({ authorization, release, server }) =>
|
({ authorization, server }) =>
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const credentials = yield* Credential.Service
|
const credentials = yield* Credential.Service
|
||||||
const catalog = yield* Catalog.Service
|
const catalog = yield* Catalog.Service
|
||||||
|
|
@ -237,15 +234,9 @@ describe("OpencodePlugin", () => {
|
||||||
})
|
})
|
||||||
|
|
||||||
yield* addPlugin()
|
yield* addPlugin()
|
||||||
expect(authorization).toEqual([])
|
expect(authorization).toEqual(["Bearer secret"])
|
||||||
release()
|
|
||||||
|
|
||||||
const provider = required(
|
const provider = required(yield* catalog.provider.get(ProviderV2.ID.make("remote")))
|
||||||
yield* eventually(
|
|
||||||
catalog.provider.get(ProviderV2.ID.make("remote")),
|
|
||||||
(item) => item?.integrationID === Integration.ID.make("opencode"),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
expect(provider).toMatchObject({
|
expect(provider).toMatchObject({
|
||||||
name: "Remote",
|
name: "Remote",
|
||||||
integrationID: "opencode",
|
integrationID: "opencode",
|
||||||
|
|
@ -283,7 +274,6 @@ describe("OpencodePlugin", () => {
|
||||||
required(yield* catalog.model.get(ProviderV2.ID.make("remote"), ModelV2.ID.make("disabled"))).enabled,
|
required(yield* catalog.model.get(ProviderV2.ID.make("remote"), ModelV2.ID.make("disabled"))).enabled,
|
||||||
).toBe(false)
|
).toBe(false)
|
||||||
expect(yield* catalog.model.get(ProviderV2.ID.make("remote"), ModelV2.ID.make("stale"))).toBeDefined()
|
expect(yield* catalog.model.get(ProviderV2.ID.make("remote"), ModelV2.ID.make("stale"))).toBeDefined()
|
||||||
expect(authorization).toContain("Bearer secret")
|
|
||||||
}),
|
}),
|
||||||
({ server }) => Effect.promise(() => server.stop(true)),
|
({ server }) => Effect.promise(() => server.stop(true)),
|
||||||
),
|
),
|
||||||
|
|
|
||||||
|
|
@ -115,6 +115,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
||||||
})
|
})
|
||||||
const messageIndex = new Map<string, Map<string, number>>()
|
const messageIndex = new Map<string, Map<string, number>>()
|
||||||
let bootstrapping: Promise<void> | undefined
|
let bootstrapping: Promise<void> | undefined
|
||||||
|
let connected = false
|
||||||
|
|
||||||
function setSessionStatus(sessionID: string, status: DataSessionStatus) {
|
function setSessionStatus(sessionID: string, status: DataSessionStatus) {
|
||||||
setStore("session", "status", sessionID, status)
|
setStore("session", "status", sessionID, status)
|
||||||
|
|
@ -845,8 +846,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
||||||
return position === undefined ? undefined : messages?.[position]
|
return position === undefined ? undefined : messages?.[position]
|
||||||
},
|
},
|
||||||
async refresh(sessionID: string) {
|
async refresh(sessionID: string) {
|
||||||
setStore("session", "message", sessionID, [])
|
|
||||||
messageIndex.set(sessionID, new Map())
|
|
||||||
const messages = mutable(
|
const messages = mutable(
|
||||||
(await sdk.api.message.list({ sessionID, limit: 200, order: "desc" })).data,
|
(await sdk.api.message.list({ sessionID, limit: 200, order: "desc" })).data,
|
||||||
).toReversed()
|
).toReversed()
|
||||||
|
|
@ -1110,8 +1109,10 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
||||||
onCleanup(
|
onCleanup(
|
||||||
sdk.event.listen(({ details }) => {
|
sdk.event.listen(({ details }) => {
|
||||||
if (details.type === "server.connected") {
|
if (details.type === "server.connected") {
|
||||||
|
const messages = connected ? Object.keys(store.session.message) : []
|
||||||
|
connected = true
|
||||||
refreshActive()
|
refreshActive()
|
||||||
void bootstrap()
|
void Promise.allSettled([bootstrap(), ...messages.map(result.session.message.refresh)])
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
handleEvent(details)
|
handleEvent(details)
|
||||||
|
|
|
||||||
|
|
@ -458,8 +458,9 @@ test("restores running manual compaction before applying live deltas", async ()
|
||||||
|
|
||||||
test("reconnects the event stream and bootstraps fresh data", async () => {
|
test("reconnects the event stream and bootstraps fresh data", async () => {
|
||||||
const events = createEventStream()
|
const events = createEventStream()
|
||||||
const requests = { active: 0, event: 0, model: 0 }
|
const requests = { active: 0, event: 0, message: 0, model: 0 }
|
||||||
let resolveActive!: (response: Response) => void
|
let resolveActive!: (response: Response) => void
|
||||||
|
let resolveMessages!: (response: Response) => void
|
||||||
const calls = createFetch((url) => {
|
const calls = createFetch((url) => {
|
||||||
if (url.pathname === "/api/event") {
|
if (url.pathname === "/api/event") {
|
||||||
requests.event++
|
requests.event++
|
||||||
|
|
@ -472,6 +473,17 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
|
||||||
resolveActive = resolve
|
resolveActive = resolve
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
if (url.pathname === "/api/session/session-stale/message") {
|
||||||
|
requests.message++
|
||||||
|
if (requests.message === 1)
|
||||||
|
return json({
|
||||||
|
data: [{ id: "message-stale", type: "user", text: "Stale", time: { created: 1 } }],
|
||||||
|
cursor: {},
|
||||||
|
})
|
||||||
|
return new Promise<Response>((resolve) => {
|
||||||
|
resolveMessages = resolve
|
||||||
|
})
|
||||||
|
}
|
||||||
if (url.pathname !== "/api/model") return
|
if (url.pathname !== "/api/model") return
|
||||||
requests.model++
|
requests.model++
|
||||||
return json({
|
return json({
|
||||||
|
|
@ -517,6 +529,8 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
|
||||||
try {
|
try {
|
||||||
await wait(() => data.location.model.list()?.[0]?.id === "model-1")
|
await wait(() => data.location.model.list()?.[0]?.id === "model-1")
|
||||||
await wait(() => data.session.status("session-stale") === "running")
|
await wait(() => data.session.status("session-stale") === "running")
|
||||||
|
await data.session.message.refresh("session-stale")
|
||||||
|
expect(data.session.message.get("session-stale", "message-stale")?.id).toBe("message-stale")
|
||||||
expect(sdk.connection.status()).toBe("connected")
|
expect(sdk.connection.status()).toBe("connected")
|
||||||
expect(sdk.connection.attempt()).toBe(0)
|
expect(sdk.connection.attempt()).toBe(0)
|
||||||
|
|
||||||
|
|
@ -530,8 +544,19 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
|
||||||
|
|
||||||
await wait(() => data.location.model.list()?.[0]?.id === "model-2", 4000)
|
await wait(() => data.location.model.list()?.[0]?.id === "model-2", 4000)
|
||||||
await wait(() => data.session.status("session-stale") === "idle")
|
await wait(() => data.session.status("session-stale") === "idle")
|
||||||
expect(data.session.status("session-new")).toBe("running")
|
await wait(() => requests.message === 2)
|
||||||
|
expect(data.session.message.get("session-stale", "message-stale")?.id).toBe("message-stale")
|
||||||
|
resolveMessages(
|
||||||
|
json({
|
||||||
|
data: [{ id: "message-fresh", type: "user", text: "Fresh", time: { created: 2 } }],
|
||||||
|
cursor: {},
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
await wait(() => data.session.message.get("session-stale", "message-fresh") !== undefined)
|
||||||
|
expect(data.session.message.get("session-stale", "message-stale")).toBeUndefined()
|
||||||
|
await wait(() => data.session.status("session-new") === "running")
|
||||||
expect(requests.event).toBe(2)
|
expect(requests.event).toBe(2)
|
||||||
|
expect(requests.message).toBe(2)
|
||||||
expect(sdk.connection.status()).toBe("connected")
|
expect(sdk.connection.status()).toBe("connected")
|
||||||
expect(sdk.connection.attempt()).toBe(0)
|
expect(sdk.connection.attempt()).toBe(0)
|
||||||
expect(sdk.connection.error()).toBeUndefined()
|
expect(sdk.connection.error()).toBeUndefined()
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue