core: ensure events publish reliably after database operations complete

This commit is contained in:
Dax Raad 2026-01-26 12:00:44 -05:00
commit 2b05833c32
2 changed files with 132 additions and 80 deletions

View file

@ -256,9 +256,17 @@ export namespace Session {
export const touch = fn(Identifier.schema("session"), async (sessionID) => { export const touch = fn(Identifier.schema("session"), async (sessionID) => {
const now = Date.now() const now = Date.now()
Database.use((db) => db.update(SessionTable).set({ time_updated: now }).where(eq(SessionTable.id, sessionID)).run()) Database.use((db) => {
const info = await get(sessionID) const row = db
Bus.publish(Event.Updated, { info }) .update(SessionTable)
.set({ time_updated: now })
.where(eq(SessionTable.id, sessionID))
.returning()
.get()
if (!row) throw new NotFoundError({ message: `Session not found: ${sessionID}` })
const info = fromRow(row)
Database.effect(() => Bus.publish(Event.Updated, { info }))
})
}) })
export async function createNext(input: { export async function createNext(input: {
@ -283,9 +291,13 @@ export namespace Session {
}, },
} }
log.info("created", result) log.info("created", result)
Database.use((db) => db.insert(SessionTable).values(toRow(result)).run()) Database.use((db) => {
Bus.publish(Event.Created, { db.insert(SessionTable).values(toRow(result)).run()
info: result, Database.effect(() =>
Bus.publish(Event.Created, {
info: result,
}),
)
}) })
const cfg = await Config.get() const cfg = await Config.get()
if (!result.parentID && (Flag.OPENCODE_AUTO_SHARE || cfg.share === "auto")) if (!result.parentID && (Flag.OPENCODE_AUTO_SHARE || cfg.share === "auto"))
@ -323,9 +335,12 @@ export namespace Session {
} }
const { ShareNext } = await import("@/share/share-next") const { ShareNext } = await import("@/share/share-next")
const share = await ShareNext.create(id) const share = await ShareNext.create(id)
Database.use((db) => db.update(SessionTable).set({ share_url: share.url }).where(eq(SessionTable.id, id)).run()) Database.use((db) => {
const info = await get(id) const row = db.update(SessionTable).set({ share_url: share.url }).where(eq(SessionTable.id, id)).returning().get()
Bus.publish(Event.Updated, { info }) if (!row) throw new NotFoundError({ message: `Session not found: ${id}` })
const info = fromRow(row)
Database.effect(() => Bus.publish(Event.Updated, { info }))
})
return share return share
}) })
@ -333,9 +348,12 @@ export namespace Session {
// Use ShareNext to remove the share (same as share function uses ShareNext to create) // Use ShareNext to remove the share (same as share function uses ShareNext to create)
const { ShareNext } = await import("@/share/share-next") const { ShareNext } = await import("@/share/share-next")
await ShareNext.remove(id) await ShareNext.remove(id)
Database.use((db) => db.update(SessionTable).set({ share_url: null }).where(eq(SessionTable.id, id)).run()) Database.use((db) => {
const info = await get(id) const row = db.update(SessionTable).set({ share_url: null }).where(eq(SessionTable.id, id)).returning().get()
Bus.publish(Event.Updated, { info }) if (!row) throw new NotFoundError({ message: `Session not found: ${id}` })
const info = fromRow(row)
Database.effect(() => Bus.publish(Event.Updated, { info }))
})
}) })
export const setTitle = fn( export const setTitle = fn(
@ -344,12 +362,18 @@ export namespace Session {
title: z.string(), title: z.string(),
}), }),
async (input) => { async (input) => {
Database.use((db) => return Database.use((db) => {
db.update(SessionTable).set({ title: input.title }).where(eq(SessionTable.id, input.sessionID)).run(), const row = db
) .update(SessionTable)
const info = await get(input.sessionID) .set({ title: input.title })
Bus.publish(Event.Updated, { info }) .where(eq(SessionTable.id, input.sessionID))
return info .returning()
.get()
if (!row) throw new NotFoundError({ message: `Session not found: ${input.sessionID}` })
const info = fromRow(row)
Database.effect(() => Bus.publish(Event.Updated, { info }))
return info
})
}, },
) )
@ -359,12 +383,18 @@ export namespace Session {
time: z.number().optional(), time: z.number().optional(),
}), }),
async (input) => { async (input) => {
Database.use((db) => return Database.use((db) => {
db.update(SessionTable).set({ time_archived: input.time }).where(eq(SessionTable.id, input.sessionID)).run(), const row = db
) .update(SessionTable)
const info = await get(input.sessionID) .set({ time_archived: input.time })
Bus.publish(Event.Updated, { info }) .where(eq(SessionTable.id, input.sessionID))
return info .returning()
.get()
if (!row) throw new NotFoundError({ message: `Session not found: ${input.sessionID}` })
const info = fromRow(row)
Database.effect(() => Bus.publish(Event.Updated, { info }))
return info
})
}, },
) )
@ -374,16 +404,18 @@ export namespace Session {
permission: PermissionNext.Ruleset, permission: PermissionNext.Ruleset,
}), }),
async (input) => { async (input) => {
Database.use((db) => return Database.use((db) => {
db const row = db
.update(SessionTable) .update(SessionTable)
.set({ permission: input.permission, time_updated: Date.now() }) .set({ permission: input.permission, time_updated: Date.now() })
.where(eq(SessionTable.id, input.sessionID)) .where(eq(SessionTable.id, input.sessionID))
.run(), .returning()
) .get()
const info = await get(input.sessionID) if (!row) throw new NotFoundError({ message: `Session not found: ${input.sessionID}` })
Bus.publish(Event.Updated, { info }) const info = fromRow(row)
return info Database.effect(() => Bus.publish(Event.Updated, { info }))
return info
})
}, },
) )
@ -394,8 +426,8 @@ export namespace Session {
summary: Info.shape.summary, summary: Info.shape.summary,
}), }),
async (input) => { async (input) => {
Database.use((db) => return Database.use((db) => {
db const row = db
.update(SessionTable) .update(SessionTable)
.set({ .set({
revert_message_id: input.revert?.messageID ?? null, revert_message_id: input.revert?.messageID ?? null,
@ -408,17 +440,19 @@ export namespace Session {
time_updated: Date.now(), time_updated: Date.now(),
}) })
.where(eq(SessionTable.id, input.sessionID)) .where(eq(SessionTable.id, input.sessionID))
.run(), .returning()
) .get()
const info = await get(input.sessionID) if (!row) throw new NotFoundError({ message: `Session not found: ${input.sessionID}` })
Bus.publish(Event.Updated, { info }) const info = fromRow(row)
return info Database.effect(() => Bus.publish(Event.Updated, { info }))
return info
})
}, },
) )
export const clearRevert = fn(Identifier.schema("session"), async (sessionID) => { export const clearRevert = fn(Identifier.schema("session"), async (sessionID) => {
Database.use((db) => return Database.use((db) => {
db const row = db
.update(SessionTable) .update(SessionTable)
.set({ .set({
revert_message_id: null, revert_message_id: null,
@ -428,11 +462,13 @@ export namespace Session {
time_updated: Date.now(), time_updated: Date.now(),
}) })
.where(eq(SessionTable.id, sessionID)) .where(eq(SessionTable.id, sessionID))
.run(), .returning()
) .get()
const info = await get(sessionID) if (!row) throw new NotFoundError({ message: `Session not found: ${sessionID}` })
Bus.publish(Event.Updated, { info }) const info = fromRow(row)
return info Database.effect(() => Bus.publish(Event.Updated, { info }))
return info
})
}) })
export const setSummary = fn( export const setSummary = fn(
@ -441,8 +477,8 @@ export namespace Session {
summary: Info.shape.summary, summary: Info.shape.summary,
}), }),
async (input) => { async (input) => {
Database.use((db) => return Database.use((db) => {
db const row = db
.update(SessionTable) .update(SessionTable)
.set({ .set({
summary_additions: input.summary?.additions, summary_additions: input.summary?.additions,
@ -451,11 +487,13 @@ export namespace Session {
time_updated: Date.now(), time_updated: Date.now(),
}) })
.where(eq(SessionTable.id, input.sessionID)) .where(eq(SessionTable.id, input.sessionID))
.run(), .returning()
) .get()
const info = await get(input.sessionID) if (!row) throw new NotFoundError({ message: `Session not found: ${input.sessionID}` })
Bus.publish(Event.Updated, { info }) const info = fromRow(row)
return info Database.effect(() => Bus.publish(Event.Updated, { info }))
return info
})
}, },
) )
@ -506,9 +544,13 @@ export namespace Session {
} }
await unshare(sessionID).catch(() => {}) await unshare(sessionID).catch(() => {})
// CASCADE delete handles messages and parts automatically // CASCADE delete handles messages and parts automatically
Database.use((db) => db.delete(SessionTable).where(eq(SessionTable.id, sessionID)).run()) Database.use((db) => {
Bus.publish(Event.Deleted, { db.delete(SessionTable).where(eq(SessionTable.id, sessionID)).run()
info: session, Database.effect(() =>
Bus.publish(Event.Deleted, {
info: session,
}),
)
}) })
} catch (e) { } catch (e) {
log.error(e) log.error(e)
@ -517,9 +559,8 @@ export namespace Session {
export const updateMessage = fn(MessageV2.Info, async (msg) => { export const updateMessage = fn(MessageV2.Info, async (msg) => {
const created_at = msg.role === "user" ? msg.time.created : msg.time.created const created_at = msg.role === "user" ? msg.time.created : msg.time.created
Database.use((db) => Database.use((db) => {
db db.insert(MessageTable)
.insert(MessageTable)
.values({ .values({
id: msg.id, id: msg.id,
session_id: msg.sessionID, session_id: msg.sessionID,
@ -527,10 +568,12 @@ export namespace Session {
data: msg, data: msg,
}) })
.onConflictDoUpdate({ target: MessageTable.id, set: { data: msg } }) .onConflictDoUpdate({ target: MessageTable.id, set: { data: msg } })
.run(), .run()
) Database.effect(() =>
Bus.publish(MessageV2.Event.Updated, { Bus.publish(MessageV2.Event.Updated, {
info: msg, info: msg,
}),
)
}) })
return msg return msg
}) })
@ -542,10 +585,14 @@ export namespace Session {
}), }),
async (input) => { async (input) => {
// CASCADE delete handles parts automatically // CASCADE delete handles parts automatically
Database.use((db) => db.delete(MessageTable).where(eq(MessageTable.id, input.messageID)).run()) Database.use((db) => {
Bus.publish(MessageV2.Event.Removed, { db.delete(MessageTable).where(eq(MessageTable.id, input.messageID)).run()
sessionID: input.sessionID, Database.effect(() =>
messageID: input.messageID, Bus.publish(MessageV2.Event.Removed, {
sessionID: input.sessionID,
messageID: input.messageID,
}),
)
}) })
return input.messageID return input.messageID
}, },
@ -558,11 +605,15 @@ export namespace Session {
partID: Identifier.schema("part"), partID: Identifier.schema("part"),
}), }),
async (input) => { async (input) => {
Database.use((db) => db.delete(PartTable).where(eq(PartTable.id, input.partID)).run()) Database.use((db) => {
Bus.publish(MessageV2.Event.PartRemoved, { db.delete(PartTable).where(eq(PartTable.id, input.partID)).run()
sessionID: input.sessionID, Database.effect(() =>
messageID: input.messageID, Bus.publish(MessageV2.Event.PartRemoved, {
partID: input.partID, sessionID: input.sessionID,
messageID: input.messageID,
partID: input.partID,
}),
)
}) })
return input.partID return input.partID
}, },
@ -583,9 +634,8 @@ export namespace Session {
export const updatePart = fn(UpdatePartInput, async (input) => { export const updatePart = fn(UpdatePartInput, async (input) => {
const part = "delta" in input ? input.part : input const part = "delta" in input ? input.part : input
const delta = "delta" in input ? input.delta : undefined const delta = "delta" in input ? input.delta : undefined
Database.use((db) => Database.use((db) => {
db db.insert(PartTable)
.insert(PartTable)
.values({ .values({
id: part.id, id: part.id,
message_id: part.messageID, message_id: part.messageID,
@ -593,11 +643,13 @@ export namespace Session {
data: part, data: part,
}) })
.onConflictDoUpdate({ target: PartTable.id, set: { data: part } }) .onConflictDoUpdate({ target: PartTable.id, set: { data: part } })
.run(), .run()
) Database.effect(() =>
Bus.publish(MessageV2.Event.PartUpdated, { Bus.publish(MessageV2.Event.PartUpdated, {
part, part,
delta, delta,
}),
)
}) })
return part return part
}) })

View file

@ -108,7 +108,7 @@ export namespace Database {
} }
} }
export function effect(fn: () => void | Promise<void>) { export function effect(fn: () => any | Promise<any>) {
try { try {
ctx.use().effects.push(fn) ctx.use().effects.push(fn)
} catch { } catch {