Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions packages/omo-senpi/changes.md
Original file line number Diff line number Diff line change
@@ -1,4 +1,16 @@
## computer use: forward the macOS canary policy (#8945)
## 2026-09-27 - Persist mailbox operations without whole-queue rewrites

The ordered-delivery mailbox now records one durable journal update per enqueue
or removal instead of serializing and fsyncing the complete pending queue after
every mutation. Enqueue and non-compacting removal updates append one event;
bounded snapshots compact drained or long journals. Existing `mailbox.json`
snapshots migrate on first open, preserving sequence numbers and queued
messages. The cap-and-restart test keeps the same count, byte, overflow, and
recovery contracts with injected limits, and the earlier 15-second timeout
override is removed.

## local launcher: `omo update` points at bun

The shipped extension now passes `computer.macos_canary` through the desktop
service to the native session. Explicit `off` reaches the macOS backend;
Expand Down
4 changes: 2 additions & 2 deletions packages/omo-senpi/plugin/extensions/omo.js

Large diffs are not rendered by default.

106 changes: 106 additions & 0 deletions packages/omo-senpi/src/components/thread/mailbox-journal-codec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,106 @@
import type { MailboxItem } from "./mailbox"

export type MailboxSnapshot = {
readonly version: 1
readonly kind: "snapshot"
readonly next_seq: number
readonly items: readonly MailboxItem[]
}

export type MailboxEnqueue = {
readonly version: 1
readonly kind: "enqueue"
readonly item: MailboxItem
}

export type MailboxRemove = {
readonly version: 1
readonly kind: "remove"
readonly message_seq: number
}

export type MailboxJournalEvent = MailboxSnapshot | MailboxEnqueue | MailboxRemove

export type LegacyStoredState = {
readonly next_seq: number
readonly queues: Readonly<Record<string, readonly MailboxItem[]>>
}

export function parseLegacyStoredState(value: unknown, path: string): LegacyStoredState {
if (
!isRecord(value) ||
!Number.isInteger(value.next_seq) ||
typeof value.next_seq !== "number" ||
value.next_seq < 1 ||
!isRecord(value.queues)
) {
throw new Error(`Invalid legacy mailbox state: ${path}`)
}
const queues: Record<string, readonly MailboxItem[]> = {}
for (const [target, queue] of Object.entries(value.queues)) {
if (!Array.isArray(queue)) throw new Error(`Invalid legacy mailbox queue: ${target}`)
queues[target] = queue.map(parseMailboxItem)
}
return { next_seq: value.next_seq, queues }
}

export function parseJournalEvent(value: unknown): MailboxJournalEvent {
if (!isRecord(value) || value.version !== 1 || typeof value.kind !== "string") {
throw new Error("Invalid mailbox journal event.")
}
if (value.kind === "enqueue") {
return { version: 1, kind: "enqueue", item: parseMailboxItem(value.item) }
}
if (value.kind === "remove" && typeof value.message_seq === "number" && Number.isInteger(value.message_seq)) {
return { version: 1, kind: "remove", message_seq: value.message_seq }
}
if (
value.kind === "snapshot" &&
typeof value.next_seq === "number" &&
Number.isInteger(value.next_seq) &&
value.next_seq >= 1 &&
Array.isArray(value.items)
) {
return {
version: 1,
kind: "snapshot",
next_seq: value.next_seq,
items: value.items.map(parseMailboxItem),
}
}
throw new Error("Invalid mailbox journal event.")
}

function parseMailboxItem(value: unknown): MailboxItem {
if (
!isRecord(value) ||
typeof value.target !== "string" ||
typeof value.message !== "string" ||
typeof value.message_seq !== "number" ||
!Number.isInteger(value.message_seq) ||
value.message_seq < 1 ||
!isDeliveryMode(value.delivery) ||
typeof value.operation_id !== "string" ||
typeof value.accepted_at !== "string" ||
(value.expected_turn_id !== undefined && typeof value.expected_turn_id !== "string")
) {
throw new Error("Invalid mailbox item.")
}
return {
target: value.target,
message: value.message,
message_seq: value.message_seq,
delivery: value.delivery,
...(value.expected_turn_id === undefined ? {} : { expected_turn_id: value.expected_turn_id }),
operation_id: value.operation_id,
accepted_at: value.accepted_at,
}
}

function isDeliveryMode(value: unknown): value is MailboxItem["delivery"] {
return value === "auto" || value === "steer" || value === "follow_up"
}

function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === "object" && value !== null && !Array.isArray(value)
}
165 changes: 165 additions & 0 deletions packages/omo-senpi/src/components/thread/mailbox-journal.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
import {
closeSync,
existsSync,
fstatSync,
fsyncSync,
ftruncateSync,
mkdirSync,
openSync,
readFileSync,
rmSync,
writeSync,
} from "node:fs"
import { join } from "node:path"

import { writeFileAtomically } from "@oh-my-opencode/utils/atomic-write"

import type { MailboxItem } from "./mailbox"
import {
type LegacyStoredState,
type MailboxJournalEvent,
type MailboxSnapshot,
parseJournalEvent,
parseLegacyStoredState,
} from "./mailbox-journal-codec"

const JOURNAL_COMPACT_EVENTS = 256

export type MailboxJournal = {
readonly nextSequence: () => number
readonly pending: (target: string) => readonly MailboxItem[]
readonly enqueue: (item: MailboxItem) => void
readonly remove: (messageSeq: number) => void
}

export function createMailboxJournal(directory: string): MailboxJournal {
mkdirSync(directory, { recursive: true, mode: 0o700 })
const journalPath = join(directory, "mailbox.jsonl")
const legacyPath = join(directory, "mailbox.json")
const loaded = loadJournal(journalPath)
const legacy = loaded === null ? loadLegacyState(legacyPath) : null
const items = new Map<number, MailboxItem>()
let nextSequence = 1
let eventCount = 0

if (loaded !== null) {
nextSequence = loaded.nextSequence
eventCount = loaded.eventCount
for (const item of loaded.items) items.set(item.message_seq, item)
rmSync(legacyPath, { force: true })
} else if (legacy !== null) {
for (const queue of Object.values(legacy.queues)) {
for (const item of queue) items.set(item.message_seq, item)
}
nextSequence = Math.max(legacy.next_seq, ...[...items.values()].map((item) => item.message_seq + 1))
writeSnapshot(journalPath, nextSequence, items.values())
eventCount = 1
rmSync(legacyPath, { force: true })
} else {
writeSnapshot(journalPath, nextSequence, [])
eventCount = 1
}

return {
nextSequence: () => nextSequence,
pending: (target) => [...items.values()]
.filter((item) => item.target === target)
.toSorted((left, right) => left.message_seq - right.message_seq),
enqueue: (item) => {
if (item.message_seq !== nextSequence) {
throw new Error(`Mailbox sequence mismatch: expected ${nextSequence}, received ${item.message_seq}.`)
}
appendEvent(journalPath, { version: 1, kind: "enqueue", item })
items.set(item.message_seq, item)
nextSequence++
eventCount++
},
remove: (messageSeq) => {
if (!items.has(messageSeq)) return
if (items.size === 1 || eventCount + 1 >= JOURNAL_COMPACT_EVENTS) {
const remaining = [...items.values()].filter((item) => item.message_seq !== messageSeq)
writeSnapshot(journalPath, nextSequence, remaining)
items.delete(messageSeq)
eventCount = 1
return
}
appendEvent(journalPath, { version: 1, kind: "remove", message_seq: messageSeq })
items.delete(messageSeq)
eventCount++
},
}
}

function loadJournal(path: string): {
readonly nextSequence: number
readonly items: readonly MailboxItem[]
readonly eventCount: number
} | null {
if (!existsSync(path)) return null
let nextSequence = 1
const items = new Map<number, MailboxItem>()
const content = readFileSync(path, "utf8")
const completeBytes = content.lastIndexOf("\n") + 1
const completeContent = content.slice(0, completeBytes)
if (completeBytes !== content.length) writeFileAtomically(path, completeContent)
const lines = completeContent.split("\n").filter((line) => line.length > 0)
if (lines.length === 0) return null
for (const line of lines) {
const event = parseJournalEvent(JSON.parse(line))
switch (event.kind) {
case "snapshot":
items.clear()
for (const item of event.items) items.set(item.message_seq, item)
nextSequence = Math.max(event.next_seq, ...event.items.map((item) => item.message_seq + 1))
break
case "enqueue":
items.set(event.item.message_seq, event.item)
nextSequence = Math.max(nextSequence, event.item.message_seq + 1)
break
case "remove":
items.delete(event.message_seq)
break
default: {
const unreachable: never = event
throw new Error(`Unexpected mailbox journal event: ${String(unreachable)}`)
}
}
}
return { nextSequence, items: [...items.values()], eventCount: lines.length }
}

function loadLegacyState(path: string): LegacyStoredState | null {
if (!existsSync(path)) return null
return parseLegacyStoredState(JSON.parse(readFileSync(path, "utf8")), path)
}

function appendEvent(path: string, event: MailboxJournalEvent): void {
const fileDescriptor = openSync(path, "r+")
const originalSize = fstatSync(fileDescriptor).size
try {
const content = Buffer.from(`${JSON.stringify(event)}\n`, "utf8")
let written = 0
while (written < content.length) {
const count = writeSync(fileDescriptor, content, written, content.length - written, originalSize + written)
if (count === 0) throw new Error("Mailbox journal append made no progress.")
written += count
}
fsyncSync(fileDescriptor)
} catch (error) {
ftruncateSync(fileDescriptor, originalSize)
fsyncSync(fileDescriptor)
throw error
} finally {
closeSync(fileDescriptor)
}
}

function writeSnapshot(path: string, nextSequence: number, items: Iterable<MailboxItem>): void {
const snapshot: MailboxSnapshot = {
version: 1,
kind: "snapshot",
next_seq: nextSequence,
items: [...items].toSorted((left, right) => left.message_seq - right.message_seq),
}
writeFileAtomically(path, `${JSON.stringify(snapshot)}\n`)
}
101 changes: 101 additions & 0 deletions packages/omo-senpi/src/components/thread/mailbox-persistence.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
import { appendFileSync, existsSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"
import { tmpdir } from "node:os"
import { join } from "node:path"

import { describe, expect, test } from "bun:test"

import { createOrderedDeliveryMailbox, type MailboxTargetPort } from "./mailbox"

const blocked: MailboxTargetPort = {
snapshot: async () => ({ turn_id: "busy", active: true }),
steer: async () => { await new Promise<void>(() => {}) },
start: async () => ({ turn_id: "new" }),
}

describe("ordered delivery mailbox persistence", () => {
test("migrates legacy whole-queue state without losing sequence or pending messages", async () => {
const directory = mkdtempSync(join(tmpdir(), "omo-mailbox-legacy-"))
const legacyItems = ["legacy-1", "legacy-2"].map((message, index) => ({
target: "target",
message,
message_seq: index + 1,
delivery: "follow_up",
operation_id: `target-${index + 1}`,
accepted_at: "2026-09-27T00:00:00.000Z",
}))
writeFileSync(join(directory, "mailbox.json"), JSON.stringify({
next_seq: 3,
queues: { target: legacyItems },
}))

const mailbox = createOrderedDeliveryMailbox({ directory, portFor: () => blocked, max_messages: 3 })
expect(mailbox.pending("target").map((item) => item.message)).toEqual(["legacy-1", "legacy-2"])
expect(existsSync(join(directory, "mailbox.json"))).toBe(false)
const accepted = await mailbox.accept("target", "new", { delivery: "follow_up" })
expect(accepted).toMatchObject({ kind: "ok", delivery: "queued", message_seq: 3, queue_position: 3 })

mailbox.close()
rmSync(directory, { recursive: true, force: true })
})

test("compacts a drained journal without reusing the last sequence after restart", async () => {
const directory = mkdtempSync(join(tmpdir(), "omo-mailbox-compaction-"))
const live: MailboxTargetPort = {
snapshot: async () => ({ turn_id: "turn-1", active: true }),
steer: async () => {},
start: async () => ({ turn_id: "new" }),
}
const first = createOrderedDeliveryMailbox({ directory, portFor: () => live })
expect(await first.accept("target", "delivered", { delivery: "steer", expected_turn_id: "turn-1" }))
.toMatchObject({ kind: "ok", delivery: "steered", message_seq: 1 })
first.close()
const restarted = createOrderedDeliveryMailbox({ directory, portFor: () => blocked })
expect(restarted.pending("target")).toHaveLength(0)
expect(await restarted.accept("target", "queued", { delivery: "follow_up" }))
.toMatchObject({ kind: "ok", delivery: "queued", message_seq: 2 })

restarted.close()
rmSync(directory, { recursive: true, force: true })
})

test("repairs a torn final journal record before accepting another message", async () => {
const directory = mkdtempSync(join(tmpdir(), "omo-mailbox-torn-tail-"))
const first = createOrderedDeliveryMailbox({ directory, portFor: () => blocked })
await first.accept("target", "first", { delivery: "follow_up" })
first.close()
appendFileSync(join(directory, "mailbox.jsonl"), '{"version":1,"kind":"enq')

const recovered = createOrderedDeliveryMailbox({ directory, portFor: () => blocked })
expect(recovered.pending("target").map((item) => item.message)).toEqual(["first"])
await recovered.accept("target", "second", { delivery: "follow_up" })
recovered.close()
const reopened = createOrderedDeliveryMailbox({ directory, portFor: () => blocked })
expect(reopened.pending("target").map((item) => item.message)).toEqual(["first", "second"])

reopened.close()
rmSync(directory, { recursive: true, force: true })
})

test("replays a removal without dropping another target's pending message", async () => {
const directory = mkdtempSync(join(tmpdir(), "omo-mailbox-remove-replay-"))
const live: MailboxTargetPort = {
snapshot: async () => ({ turn_id: "turn-1", active: true }),
steer: async () => {},
start: async () => ({ turn_id: "new" }),
}
const mailbox = createOrderedDeliveryMailbox({
directory,
portFor: (target) => target === "pending" ? blocked : live,
})
await mailbox.accept("pending", "keep", { delivery: "follow_up" })
await mailbox.accept("delivered", "remove", { delivery: "steer", expected_turn_id: "turn-1" })
mailbox.close()

const restarted = createOrderedDeliveryMailbox({ directory, portFor: () => blocked })
expect(restarted.pending("pending").map((item) => item.message)).toEqual(["keep"])
expect(restarted.pending("delivered")).toHaveLength(0)

restarted.close()
rmSync(directory, { recursive: true, force: true })
})
})
Loading
Loading