Compare commits

...
1 Commits
Author SHA1 Message Date
James Long 7be5a81a7b feat(core): add persistent terminal groups 2026-08-25 01:05:38 +00:00
8 changed files with 274 additions and 0 deletions
+1
View File
@@ -0,0 +1 @@
export { Group } from "./persistent-pty/group.js"
+103
View File
@@ -0,0 +1,103 @@
export * as Group from "./group.js"
import { Group } from "@opencode-ai/schema/group"
import { makeGlobalNode } from "@opencode-ai/util/effect/app-node"
import { Context, Effect, Layer, Schema, Semaphore } from "effect"
import { Bus } from "../bus.js"
import { KV } from "../kv.js"
export const ID = Group.ID
export type ID = Group.ID
export const Item = Group.Item
export type Item = Group.Item
export const Info = Group.Info
export type Info = Group.Info
export const Event = Group.Event
export interface Interface {
readonly list: () => Effect.Effect<ReadonlyArray<Info>>
readonly get: (id: ID) => Effect.Effect<Info | undefined>
readonly create: (items?: ReadonlyArray<Item>) => Effect.Effect<Info>
readonly set: (group: Info) => Effect.Effect<void>
readonly remove: (id: ID) => Effect.Effect<void>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/Group") {}
const key = "group:v1"
const Document = Schema.Array(Info)
const layer = Layer.effect(
Service,
Effect.gen(function* () {
const kv = yield* KV.Service
const bus = yield* Bus.Service
const lock = Semaphore.makeUnsafe(1)
const list = Effect.fn("Group.list")(function* () {
const value = yield* kv.get(key)
return Schema.is(Document)(value) ? value : []
})
return Service.of({
list,
get: Effect.fn("Group.get")(function* (id) {
return (yield* list()).find((group) => group.id === id)
}),
create: Effect.fn("Group.create")(function* (items = []) {
return yield* lock.withPermit(
Effect.gen(function* () {
const group = Info.make({ id: ID.create(), items: Array.from(items) })
yield* kv.set(key, (yield* list()).concat(group))
return group
}),
)
}),
set: Effect.fn("Group.set")(function* (group) {
yield* lock.withPermit(
Effect.gen(function* () {
const groups = yield* list()
const index = groups.findIndex((item) => item.id === group.id)
yield* kv.set(
key,
index === -1 ? groups.concat(group) : groups.map((item) => (item.id === group.id ? group : item)),
)
const previous = groups[index]
if (!previous) return
yield* Effect.forEach(
group.items.filter(
(item) => !previous.items.some((current) => current.type === item.type && current.id === item.id),
),
(item) => bus.publish(Event.ItemAdded, { groupID: group.id, item }),
{ discard: true },
)
yield* Effect.forEach(
previous.items.filter(
(item) => !group.items.some((next) => next.type === item.type && next.id === item.id),
),
(item) => bus.publish(Event.ItemRemoved, { groupID: group.id, item }),
{ discard: true },
)
}),
)
}),
remove: Effect.fn("Group.remove")(function* (id) {
yield* lock.withPermit(
Effect.gen(function* () {
const groups = yield* list()
const group = groups.find((group) => group.id === id)
yield* kv.set(key, groups.filter((group) => group.id !== id))
if (!group) return
yield* Effect.forEach(
group.items,
(item) => bus.publish(Event.ItemRemoved, { groupID: id, item }),
{ discard: true },
)
}),
)
}),
})
}),
)
export const node = makeGlobalNode({ service: Service, layer, deps: [KV.node, Bus.node] })
+91
View File
@@ -0,0 +1,91 @@
import { describe, expect } from "bun:test"
import { Group } from "@opencode-ai/core/persistent-pty"
import { Bus } from "@opencode-ai/core/bus"
import { KV } from "@opencode-ai/core/kv"
import { Pty } from "@opencode-ai/schema/pty"
import { Session } from "@opencode-ai/schema/session"
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
import { Effect, Fiber, Stream } from "effect"
import { testEffect } from "./lib/effect"
const it = testEffect(LayerNode.compile(LayerNode.group([Group.node, KV.node, Bus.node])))
describe("Group", () => {
it.effect("persists ordered groups in one versioned KV document", () =>
Effect.gen(function* () {
const groups = yield* Group.Service
const kv = yield* KV.Service
const created = yield* groups.create([
{ type: "session", id: Session.ID.make("ses_one") },
{ type: "terminal", id: Pty.ID.make("pty_one") },
])
expect(yield* groups.get(created.id)).toEqual(created)
expect(yield* groups.list()).toEqual([created])
expect(yield* kv.get("group:v1")).toEqual([created])
const updated = Group.Info.make({
id: created.id,
items: [{ type: "terminal", id: Pty.ID.make("pty_two") }],
})
yield* groups.set(updated)
expect(yield* groups.list()).toEqual([updated])
yield* groups.remove(created.id)
expect(yield* groups.get(created.id)).toBeUndefined()
expect(yield* kv.get("group:v1")).toEqual([])
}),
)
it.effect("serializes concurrent document mutations", () =>
Effect.gen(function* () {
const groups = yield* Group.Service
yield* Effect.all(
Array.from({ length: 20 }, (_, index) =>
groups.create([{ type: "session", id: Session.ID.make(`ses_${index}`) }]),
),
{ concurrency: "unbounded" },
)
expect(yield* groups.list()).toHaveLength(20)
}),
)
it.effect("publishes every removed group item", () =>
Effect.gen(function* () {
const groups = yield* Group.Service
const bus = yield* Bus.Service
const session = { type: "session" as const, id: Session.ID.make("ses_one") }
const terminal = { type: "terminal" as const, id: Pty.ID.make("pty_one") }
const group = yield* groups.create([session, terminal])
const events = yield* bus
.subscribe(Group.Event.ItemRemoved)
.pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
yield* Effect.yieldNow
yield* groups.set(Group.Info.make({ id: group.id, items: [session] }))
yield* groups.remove(group.id)
expect(Array.from(yield* Fiber.join(events)).map((event) => event.data)).toEqual([
{ groupID: group.id, item: terminal },
{ groupID: group.id, item: session },
])
}),
)
it.effect("publishes every added group item", () =>
Effect.gen(function* () {
const groups = yield* Group.Service
const bus = yield* Bus.Service
const session = { type: "session" as const, id: Session.ID.make("ses_one") }
const terminal = { type: "terminal" as const, id: Pty.ID.make("pty_one") }
const group = yield* groups.create([session])
const event = yield* bus.subscribe(Group.Event.ItemAdded).pipe(Stream.runHead, Effect.forkScoped)
yield* Effect.yieldNow
yield* groups.set(Group.Info.make({ id: group.id, items: [session, terminal] }))
expect((yield* Fiber.join(event)).valueOrUndefined?.data).toEqual({ groupID: group.id, item: terminal })
}),
)
})
+2
View File
@@ -10,6 +10,7 @@ import { Event } from "./event.js"
import { FileSystem } from "./filesystem.js"
import { FileSystemV1 } from "./filesystem-v1.js"
import { Form } from "./form.js"
import { Group } from "./group.js"
import { InstallationEvent } from "./installation-event.js"
import { Integration } from "./integration.js"
import { LegacyEventV1 } from "./legacy-event.js"
@@ -56,6 +57,7 @@ const featureDefinitions = Event.inventory(
...Pty.Event.Definitions,
...Shell.Event.Definitions,
...Form.Event.Definitions,
...Group.Event.Definitions,
...WebSearch.Event.Definitions,
)
+43
View File
@@ -0,0 +1,43 @@
export * as Group from "./group.js"
import { Schema } from "effect"
import { ephemeral, inventory } from "./event.js"
import { ascending } from "./identifier.js"
import { Pty } from "./pty.js"
import { statics } from "./schema.js"
import { Session } from "./session.js"
const IDSchema = Schema.String.check(Schema.isStartsWith("grp_")).pipe(Schema.brand("GroupID"))
export const ID = IDSchema.pipe(
statics((schema: typeof IDSchema) => ({ create: () => schema.make("grp_" + ascending()) })),
)
export type ID = typeof ID.Type
export const SessionItem = Schema.Struct({
type: Schema.tag("session"),
id: Session.ID,
})
export interface SessionItem extends Schema.Schema.Type<typeof SessionItem> {}
export const TerminalItem = Schema.Struct({
type: Schema.tag("terminal"),
id: Pty.ID,
})
export interface TerminalItem extends Schema.Schema.Type<typeof TerminalItem> {}
export const Item = Schema.Union([SessionItem, TerminalItem]).pipe(
Schema.toTaggedUnion("type"),
Schema.annotate({ identifier: "Group.Item" }),
)
export type Item = typeof Item.Type
export const Info = Schema.Struct({
id: ID,
items: Schema.Array(Item),
}).annotate({ identifier: "Group.Info" })
export interface Info extends Schema.Schema.Type<typeof Info> {}
const ItemAdded = ephemeral({ type: "group.item.added", schema: { groupID: ID, item: Item } })
const ItemRemoved = ephemeral({ type: "group.item.removed", schema: { groupID: ID, item: Item } })
export const Event = { ItemAdded, ItemRemoved, Definitions: inventory(ItemAdded, ItemRemoved) }
+1
View File
@@ -6,6 +6,7 @@ export { Credential } from "./credential.js"
export { Event } from "./event.js"
export { FileSystem } from "./filesystem.js"
export { Form } from "./form.js"
export { Group } from "./group.js"
export { Integration } from "./integration.js"
export { LLM } from "./llm.js"
export { Location } from "./location.js"
@@ -4,6 +4,7 @@ import {
Config,
FileSystem,
Form,
Group,
Integration,
Permission,
Project,
@@ -64,6 +65,7 @@ describe("public event manifest", () => {
expect(Integration.Event.Definitions).toEqual([Integration.Event.Updated, Integration.Event.ConnectionUpdated])
expect(Permission.Event.Definitions).toEqual([Permission.Event.Asked, Permission.Event.Replied])
expect(Form.Event.Definitions).toEqual([Form.Event.Created, Form.Event.Replied, Form.Event.Cancelled])
expect(Group.Event.Definitions).toEqual([Group.Event.ItemAdded, Group.Event.ItemRemoved])
expect(Reference.Event.Definitions).toEqual([Reference.Event.Updated])
expect(Plugin.Event.Definitions).toEqual([Plugin.Event.Added, Plugin.Event.Updated])
expect(McpEvent.Definitions).toEqual([McpEvent.ToolsChanged, McpEvent.ResourcesChanged, McpEvent.StatusChanged])
+31
View File
@@ -0,0 +1,31 @@
import { describe, expect, test } from "bun:test"
import { Schema } from "effect"
import { Group } from "../src/group.js"
import { Pty } from "../src/pty.js"
import { Session } from "../src/session.js"
describe("Group", () => {
test("creates branded group IDs", () => {
expect(Group.ID.create()).toStartWith("grp_")
expect(() => Schema.decodeUnknownSync(Group.ID)("ses_invalid")).toThrow()
})
test("preserves one ordered session and terminal item list", () => {
const group = Schema.decodeUnknownSync(Group.Info)({
id: Group.ID.create(),
items: [
{ type: "session", id: Session.ID.make("ses_one") },
{ type: "terminal", id: Pty.ID.make("pty_one") },
{ type: "session", id: Session.ID.make("ses_two") },
],
})
expect(group.items.map((item) => item.type)).toEqual(["session", "terminal", "session"])
expect(() =>
Schema.decodeUnknownSync(Group.Info)({
id: group.id,
items: [{ type: "other", id: "other_one" }],
}),
).toThrow()
})
})