Skip to content

Commit 2bb15e0

Browse files
fix(session): avoid full history hydration after compaction
Factor the compaction filter into an incremental scan and have filterCompactedEffect page backward through messages, stopping as soon as a completed compaction boundary (and any retained tail) is found, instead of hydrating the entire session history first. Ported from anomalyco#31638 (already Effect-based; applied as-is). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent e212a67 commit 2bb15e0

2 files changed

Lines changed: 258 additions & 23 deletions

File tree

packages/opencode/src/session/message-v2.ts

Lines changed: 66 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -518,30 +518,45 @@ export const get = Effect.fn("MessageV2.get")(function* (input: { sessionID: Ses
518518
}
519519
})
520520

521-
export function filterCompacted(msgs: Iterable<WithParts>) {
522-
const result = [] as WithParts[]
523-
const completed = new Set<string>()
524-
let retain: MessageID | undefined
525-
for (const msg of msgs) {
526-
result.push(msg)
527-
if (retain) {
528-
if (msg.info.id === retain) break
529-
continue
530-
}
531-
if (msg.info.role === "user" && completed.has(msg.info.id)) {
532-
const part = msg.parts.find((item): item is CompactionPart => item.type === "compaction")
533-
if (!part) continue
534-
if (!part.tail_start_id) break
535-
retain = part.tail_start_id
536-
if (msg.info.id === retain) break
537-
continue
521+
type CompactedScan = {
522+
result: WithParts[]
523+
completed: Set<string>
524+
retain: MessageID | undefined
525+
done: boolean
526+
}
527+
528+
function compactedScan() {
529+
return {
530+
result: [] as WithParts[],
531+
completed: new Set<string>(),
532+
retain: undefined,
533+
done: false,
534+
} satisfies CompactedScan
535+
}
536+
537+
function scanCompacted(scan: CompactedScan, msg: WithParts) {
538+
scan.result.push(msg)
539+
if (scan.retain) {
540+
if (msg.info.id === scan.retain) scan.done = true
541+
return
542+
}
543+
if (msg.info.role === "user" && scan.completed.has(msg.info.id)) {
544+
const part = msg.parts.find((item): item is CompactionPart => item.type === "compaction")
545+
if (!part) return
546+
if (!part.tail_start_id) {
547+
scan.done = true
548+
return
538549
}
539-
if (msg.info.role === "user" && completed.has(msg.info.id) && msg.parts.some((part) => part.type === "compaction"))
540-
break
541-
if (msg.info.role === "assistant" && msg.info.summary && msg.info.finish && !msg.info.error)
542-
completed.add(msg.info.parentID)
550+
scan.retain = part.tail_start_id
551+
if (msg.info.id === scan.retain) scan.done = true
552+
return
543553
}
544-
result.reverse()
554+
if (msg.info.role === "assistant" && msg.info.summary && msg.info.finish && !msg.info.error)
555+
scan.completed.add(msg.info.parentID)
556+
}
557+
558+
function finishCompactedScan(scan: CompactedScan) {
559+
const result = scan.result.reverse()
545560
const compactionIndex = result.findLastIndex(
546561
(msg) =>
547562
msg.info.role === "user" &&
@@ -571,8 +586,36 @@ export function filterCompacted(msgs: Iterable<WithParts>) {
571586
return result
572587
}
573588

589+
export function filterCompacted(msgs: Iterable<WithParts>) {
590+
const scan = compactedScan()
591+
for (const msg of msgs) {
592+
scanCompacted(scan, msg)
593+
if (scan.done) break
594+
}
595+
return finishCompactedScan(scan)
596+
}
597+
574598
export const filterCompactedEffect = Effect.fnUntraced(function* (sessionID: SessionID) {
575-
return filterCompacted(yield* stream(sessionID))
599+
const scan = compactedScan()
600+
const size = 50
601+
let before: string | undefined
602+
while (!scan.done) {
603+
const next = yield* page({ sessionID, limit: size, before }).pipe(
604+
Effect.catchIf(NotFoundError.isInstance, () =>
605+
Effect.succeed({ items: [] as WithParts[], more: false, cursor: undefined }),
606+
),
607+
)
608+
if (next.items.length === 0) break
609+
for (let i = next.items.length - 1; i >= 0; i--) {
610+
const item = next.items[i]
611+
if (!item) continue
612+
scanCompacted(scan, item)
613+
if (scan.done) break
614+
}
615+
if (scan.done || !next.more || !next.cursor) break
616+
before = next.cursor
617+
}
618+
return finishCompactedScan(scan)
576619
})
577620

578621
// filterCompacted reorders messages for model consumption

packages/opencode/test/session/messages-pagination.test.ts

Lines changed: 192 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -700,6 +700,190 @@ describe("MessageV2.filterCompacted", () => {
700700
),
701701
)
702702

703+
it.instance("filterCompactedEffect matches full stream filtering with tail retention", () =>
704+
withSession(({ session, sessionID }) =>
705+
Effect.gen(function* () {
706+
const old = yield* fill(sessionID, 75, (i: number) => Date.now() - 100_000 + i)
707+
708+
const u1 = yield* addUser(sessionID, "first")
709+
const a1 = yield* addAssistant(sessionID, u1, { finish: "end_turn" })
710+
yield* session.updatePart({
711+
id: PartID.ascending(),
712+
sessionID,
713+
messageID: a1,
714+
type: "text",
715+
text: "first reply",
716+
})
717+
718+
const u2 = yield* addUser(sessionID, "second")
719+
const a2 = yield* addAssistant(sessionID, u2, { finish: "end_turn" })
720+
yield* session.updatePart({
721+
id: PartID.ascending(),
722+
sessionID,
723+
messageID: a2,
724+
type: "text",
725+
text: "second reply",
726+
})
727+
728+
const c1 = yield* addUser(sessionID)
729+
yield* addCompactionPart(sessionID, c1, u2)
730+
const s1 = yield* addAssistant(sessionID, c1, { summary: true, finish: "end_turn" })
731+
yield* session.updatePart({
732+
id: PartID.ascending(),
733+
sessionID,
734+
messageID: s1,
735+
type: "text",
736+
text: "summary",
737+
})
738+
739+
const u3 = yield* addUser(sessionID, "third")
740+
const a3 = yield* addAssistant(sessionID, u3, { finish: "end_turn" })
741+
yield* session.updatePart({
742+
id: PartID.ascending(),
743+
sessionID,
744+
messageID: a3,
745+
type: "text",
746+
text: "third reply",
747+
})
748+
749+
const full = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
750+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
751+
752+
expect(optimized.map((item) => item.info.id)).toEqual(full.map((item) => item.info.id))
753+
expect(optimized.map((item) => item.info.id)).toEqual([c1, s1, u2, a2, u3, a3])
754+
expect(optimized.some((item) => old.includes(item.info.id))).toBe(false)
755+
}),
756+
),
757+
)
758+
759+
it.instance("filterCompactedEffect keeps retained tail across page boundaries", () =>
760+
withSession(({ session, sessionID }) =>
761+
Effect.gen(function* () {
762+
const old = yield* fill(sessionID, 10, (i: number) => Date.now() - 200_000 + i)
763+
const tail = yield* fill(sessionID, 60, (i: number) => Date.now() - 100_000 + i)
764+
const tailStart = tail[0]
765+
if (!tailStart) throw new Error("expected retained tail")
766+
const tailEnd = tail[tail.length - 1]
767+
if (!tailEnd) throw new Error("expected retained tail end")
768+
769+
const c1 = yield* addUser(sessionID)
770+
yield* addCompactionPart(sessionID, c1, tailStart)
771+
const s1 = yield* addAssistant(sessionID, c1, { summary: true, finish: "end_turn" })
772+
yield* session.updatePart({
773+
id: PartID.ascending(),
774+
sessionID,
775+
messageID: s1,
776+
type: "text",
777+
text: "summary",
778+
})
779+
780+
const u1 = yield* addUser(sessionID, "next")
781+
const a1 = yield* addAssistant(sessionID, u1, { finish: "end_turn" })
782+
yield* session.updatePart({
783+
id: PartID.ascending(),
784+
sessionID,
785+
messageID: a1,
786+
type: "text",
787+
text: "next reply",
788+
})
789+
790+
const full = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
791+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
792+
const ids = optimized.map((item) => item.info.id)
793+
794+
expect(ids).toEqual(full.map((item) => item.info.id))
795+
expect(ids.slice(0, 2)).toEqual([c1, s1])
796+
expect(ids).toContain(tailStart)
797+
expect(ids).toContain(tailEnd)
798+
expect(ids.slice(-2)).toEqual([u1, a1])
799+
expect(ids.some((item) => old.includes(item))).toBe(false)
800+
}),
801+
),
802+
)
803+
804+
it.instance("filterCompactedEffect stops at completed compaction without retained tail", () =>
805+
withSession(({ session, sessionID }) =>
806+
Effect.gen(function* () {
807+
const old = yield* fill(sessionID, 75, (i: number) => Date.now() - 100_000 + i)
808+
809+
const c1 = yield* addUser(sessionID)
810+
yield* addCompactionPart(sessionID, c1)
811+
const s1 = yield* addAssistant(sessionID, c1, { summary: true, finish: "end_turn" })
812+
yield* session.updatePart({
813+
id: PartID.ascending(),
814+
sessionID,
815+
messageID: s1,
816+
type: "text",
817+
text: "summary",
818+
})
819+
820+
const u1 = yield* addUser(sessionID, "next")
821+
const a1 = yield* addAssistant(sessionID, u1, { finish: "end_turn" })
822+
yield* session.updatePart({
823+
id: PartID.ascending(),
824+
sessionID,
825+
messageID: a1,
826+
type: "text",
827+
text: "next reply",
828+
})
829+
830+
const full = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
831+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
832+
833+
expect(optimized.map((item) => item.info.id)).toEqual(full.map((item) => item.info.id))
834+
expect(optimized.map((item) => item.info.id)).toEqual([c1, s1, u1, a1])
835+
expect(optimized.some((item) => old.includes(item.info.id))).toBe(false)
836+
}),
837+
),
838+
)
839+
840+
it.instance("filterCompactedEffect keeps full history without compaction", () =>
841+
withSession(({ sessionID }) =>
842+
Effect.gen(function* () {
843+
yield* fill(sessionID, 75, (i: number) => Date.now() - 100_000 + i)
844+
845+
const full = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
846+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
847+
848+
expect(optimized.map((item) => item.info.id)).toEqual(full.map((item) => item.info.id))
849+
expect(optimized).toHaveLength(full.length)
850+
}),
851+
),
852+
)
853+
854+
it.instance("filterCompactedEffect keeps full history without completed compaction", () =>
855+
withSession(({ sessionID }) =>
856+
Effect.gen(function* () {
857+
const old = yield* fill(sessionID, 75, (i: number) => Date.now() - 100_000 + i)
858+
859+
const missing = yield* addUser(sessionID)
860+
yield* addCompactionPart(sessionID, missing)
861+
862+
const errorParent = yield* addUser(sessionID)
863+
yield* addCompactionPart(sessionID, errorParent)
864+
const error = new SessionV1.APIError({
865+
message: "boom",
866+
isRetryable: true,
867+
}).toObject() as SessionV1.Assistant["error"]
868+
yield* addAssistant(sessionID, errorParent, { summary: true, finish: "end_turn", error })
869+
870+
const unfinished = yield* addUser(sessionID)
871+
yield* addCompactionPart(sessionID, unfinished)
872+
yield* addAssistant(sessionID, unfinished, { summary: true })
873+
874+
const full = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
875+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
876+
const ids = optimized.map((item) => item.info.id)
877+
878+
expect(ids).toEqual(full.map((item) => item.info.id))
879+
expect(ids).toContain(missing)
880+
expect(ids).toContain(errorParent)
881+
expect(ids).toContain(unfinished)
882+
expect(ids).toContain(old[0])
883+
}),
884+
),
885+
)
886+
703887
it.instance("retains original tail when compaction stores tail_start_id", () =>
704888
withSession(({ session, sessionID }) =>
705889
Effect.gen(function* () {
@@ -745,8 +929,10 @@ describe("MessageV2.filterCompacted", () => {
745929
})
746930

747931
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
932+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
748933

749934
expect(result.map((item) => item.info.id)).toEqual([c1, s1, u2, a2, u3, a3])
935+
expect(optimized.map((item) => item.info.id)).toEqual(result.map((item) => item.info.id))
750936
}),
751937
),
752938
)
@@ -802,7 +988,9 @@ describe("MessageV2.filterCompacted", () => {
802988

803989
const forked = yield* session.fork({ sessionID: created.id })
804990
const childFiltered = MessageV2.filterCompacted(yield* MessageV2.stream(forked.id))
991+
const childOptimized = yield* MessageV2.filterCompactedEffect(forked.id)
805992
expect(childFiltered).toHaveLength(parentFiltered.length)
993+
expect(childOptimized.map((item) => item.info.id)).toEqual(childFiltered.map((item) => item.info.id))
806994

807995
const tailPart = childFiltered.flatMap((m) => m.parts).find((p) => p.type === "compaction")
808996
expect(tailPart?.type).toBe("compaction")
@@ -868,8 +1056,10 @@ describe("MessageV2.filterCompacted", () => {
8681056
})
8691057

8701058
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
1059+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
8711060

8721061
expect(result.map((item) => item.info.id)).toEqual([c1, s1, a3, u3, a4])
1062+
expect(optimized.map((item) => item.info.id)).toEqual(result.map((item) => item.info.id))
8731063
}),
8741064
),
8751065
)
@@ -940,8 +1130,10 @@ describe("MessageV2.filterCompacted", () => {
9401130
})
9411131

9421132
const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
1133+
const optimized = yield* MessageV2.filterCompactedEffect(sessionID)
9431134

9441135
expect(result.map((item) => item.info.id)).toEqual([c2, s2, u3, a3, u4, a4])
1136+
expect(optimized.map((item) => item.info.id)).toEqual(result.map((item) => item.info.id))
9451137
}),
9461138
),
9471139
)

0 commit comments

Comments
 (0)