merge: integrate origin/main (406 commits) into live branch

Merge upstream's paired-room stabilization, dedicated message-runtime-loop,
outbound delivery refactor (ipc-outbound-delivery / deliverFormattedCanonicalMessage),
discord module split, dashboard rework, and migrations 015-018 into our branch.

Conflict policy: prefer upstream for the heavily-refactored paired-room /
message-runtime / db cluster (it supersedes our earlier loop band-aids), keep
only orthogonal unique fixes:
- codex bundled-JS launcher fix (codex-warmup.ts)
- slot-aligned Claude+Codex usage primer (usage-primer.ts) re-wired into index.ts
- MemoryHigh=3G cgroup bound, discord perms diagnostic, prompt docs
- progress null-race + silent-failure publish guards (message-turn-controller.ts)

Dropped (superseded by upstream or unwired): reviewer STEP_DONE/TASK_DONE loop
band-aids, router markdown-escape override, register tribunal-default+git-init,
restart-context rewindToSeq (unwired).

Verified: typecheck clean, build clean, vitest shows zero new regressions vs the
origin/main baseline (only pre-existing env-config failures remain).
This commit is contained in:
Codex
2026-06-08 21:15:39 +09:00
409 changed files with 53523 additions and 10332 deletions

View File

@@ -1,5 +1,25 @@
import { Database } from 'bun:sqlite';
const ROOM_SKILL_OVERRIDES_SCHEMA = `
CREATE TABLE IF NOT EXISTS room_skill_overrides (
chat_jid TEXT NOT NULL,
agent_type TEXT NOT NULL,
skill_scope TEXT NOT NULL,
skill_name TEXT NOT NULL,
enabled INTEGER NOT NULL,
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (chat_jid, agent_type, skill_scope, skill_name),
FOREIGN KEY (chat_jid) REFERENCES room_settings(chat_jid) ON DELETE CASCADE,
CHECK (agent_type IN ('claude-code', 'codex')),
CHECK (enabled IN (0, 1)),
CHECK (length(skill_scope) > 0),
CHECK (length(skill_name) > 0)
);
CREATE INDEX IF NOT EXISTS idx_room_skill_overrides_room
ON room_skill_overrides(chat_jid, agent_type);
`;
export function applyBaseSchema(database: Database): void {
database.exec(`
CREATE TABLE IF NOT EXISTS chats (
@@ -24,6 +44,7 @@ export function applyBaseSchema(database: Database): void {
FOREIGN KEY (chat_jid) REFERENCES chats(jid)
);
CREATE INDEX IF NOT EXISTS idx_timestamp ON messages(timestamp);
CREATE INDEX IF NOT EXISTS idx_messages_chat_timestamp ON messages(chat_jid, timestamp DESC);
CREATE TABLE IF NOT EXISTS message_sequence (
id INTEGER PRIMARY KEY AUTOINCREMENT
);
@@ -160,7 +181,7 @@ export function applyBaseSchema(database: Database): void {
task_id TEXT NOT NULL,
turn_number INTEGER NOT NULL,
role TEXT NOT NULL,
output_text TEXT NOT NULL,
output_text TEXT NOT NULL, attachment_payload TEXT,
verdict TEXT,
created_at TEXT NOT NULL,
UNIQUE(task_id, turn_number, role)
@@ -410,8 +431,7 @@ export function applyBaseSchema(database: Database): void {
last_used_at TEXT,
archived_at TEXT
);
CREATE INDEX IF NOT EXISTS idx_memories_scope
ON memories(scope_kind, scope_key);
CREATE INDEX IF NOT EXISTS idx_memories_scope ON memories(scope_kind, scope_key);
CREATE INDEX IF NOT EXISTS idx_memories_active
ON memories(scope_kind, scope_key, archived_at, created_at);
CREATE VIRTUAL TABLE IF NOT EXISTS memories_fts USING fts5(
@@ -435,4 +455,5 @@ export function applyBaseSchema(database: Database): void {
VALUES (new.id, new.content, new.keywords_json);
END;
`);
database.exec(ROOM_SKILL_OVERRIDES_SCHEMA);
}

View File

@@ -39,7 +39,10 @@ function getExpectedSchemaMigrations(): Array<{
{ version: 12, name: 'paired_verdict_and_step_telemetry' },
{ version: 13, name: 'message_source_kind' },
{ version: 14, name: 'work_item_attachments' },
{ version: 15, name: 'reviewer_failure_count' },
{ version: 15, name: 'turn_progress_text' },
{ version: 16, name: 'room_skill_overrides' },
{ version: 17, name: 'scheduled_task_room_role' },
{ version: 18, name: 'paired_turn_output_attachments' },
];
}

View File

@@ -9,7 +9,11 @@ import {
inferRoleFromServiceShadow,
resolveRoleServiceShadow,
} from '../role-service-shadow.js';
import type { AgentType, PairedRoomRole } from '../types.js';
import {
normalizePairedRoomRoleOrNull,
type AgentType,
type PairedRoomRole,
} from '../types.js';
import { normalizeStoredAgentType } from './room-registration.js';
interface RequiredRoleMetadataInput {
@@ -20,7 +24,7 @@ interface RequiredRoleMetadataInput {
fallbackAgentType?: AgentType | null;
}
interface OptionalRoleMetadataInput extends RequiredRoleMetadataInput {}
type OptionalRoleMetadataInput = RequiredRoleMetadataInput;
function resolveFilledRequiredAgentType(
input: RequiredRoleMetadataInput,
@@ -336,14 +340,6 @@ export function readCanonicalChannelOwnerLeaseMetadata(input: {
};
}
function normalizeHandoffRole(
role: string | null | undefined,
): PairedRoomRole | null {
return role === 'owner' || role === 'reviewer' || role === 'arbiter'
? role
: null;
}
function assertStoredRoleMatchesServiceShadow(args: {
context: string;
storedRole: PairedRoomRole | null;
@@ -389,11 +385,11 @@ export function fillCanonicalServiceHandoffMetadata(input: {
}): CanonicalServiceHandoffMetadata {
const context = `service_handoffs(${input.id})`;
const sourceRole =
normalizeHandoffRole(input.source_role) ??
normalizeHandoffRole(input.intended_role);
normalizePairedRoomRoleOrNull(input.source_role) ??
normalizePairedRoomRoleOrNull(input.intended_role);
const targetRole =
normalizeHandoffRole(input.target_role) ??
normalizeHandoffRole(input.intended_role);
normalizePairedRoomRoleOrNull(input.target_role) ??
normalizePairedRoomRoleOrNull(input.intended_role);
const storedSourceAgentType = normalizeStoredAgentType(
input.source_agent_type,
);
@@ -541,11 +537,11 @@ export function readCanonicalServiceHandoffMetadata(input: {
}): CanonicalServiceHandoffMetadata {
const context = `service_handoffs(${input.id})`;
const sourceRole =
normalizeHandoffRole(input.source_role) ??
normalizeHandoffRole(input.intended_role);
normalizePairedRoomRoleOrNull(input.source_role) ??
normalizePairedRoomRoleOrNull(input.intended_role);
const targetRole =
normalizeHandoffRole(input.target_role) ??
normalizeHandoffRole(input.intended_role);
normalizePairedRoomRoleOrNull(input.target_role) ??
normalizePairedRoomRoleOrNull(input.intended_role);
const sourceServiceId = readStoredHandoffServiceId({
context,
field: 'source_service_id',

View File

@@ -95,6 +95,17 @@ export function getAllChatsFromDatabase(database: Database): ChatInfo[] {
.all() as ChatInfo[];
}
export function hasMessageInDatabase(
database: Database,
chatJid: string,
id: string,
): boolean {
const row = database
.prepare('SELECT 1 FROM messages WHERE chat_jid = ? AND id = ? LIMIT 1')
.get(chatJid, id);
return !!row;
}
export function storeMessageInDatabase(
database: Database,
msg: NewMessage,
@@ -308,22 +319,31 @@ export function getMessagesSinceSeqFromDatabase(
return rows.map(normalizeMessageRow);
}
const recentChatMessagesStmtCache = new WeakMap<
Database,
ReturnType<Database['prepare']>
>();
export function getRecentChatMessagesFromDatabase(
database: Database,
chatJid: string,
limit: number = 20,
): NewMessage[] {
const sql = `
SELECT * FROM (
SELECT id, chat_jid, sender, sender_name, content, timestamp, is_from_me, is_bot_message, message_source_kind
FROM messages
WHERE chat_jid = ?
AND content != '' AND content IS NOT NULL
ORDER BY timestamp DESC
LIMIT ?
) ORDER BY timestamp
`;
const rows = database.prepare(sql).all(chatJid, limit) as Array<
let stmt = recentChatMessagesStmtCache.get(database);
if (!stmt) {
stmt = database.prepare(`
SELECT * FROM (
SELECT id, chat_jid, sender, sender_name, content, timestamp, is_from_me, is_bot_message, message_source_kind
FROM messages
WHERE chat_jid = ?
AND content != '' AND content IS NOT NULL
ORDER BY timestamp DESC
LIMIT ?
) ORDER BY timestamp
`);
recentChatMessagesStmtCache.set(database, stmt);
}
const rows = stmt.all(chatJid, limit) as Array<
NewMessage & {
is_from_me?: boolean | number;
is_bot_message?: boolean | number;
@@ -332,36 +352,42 @@ export function getRecentChatMessagesFromDatabase(
return rows.map(normalizeMessageRow);
}
/**
* Returns the seq of the earliest unanswered human message in this chat —
* i.e. the first user message that came after the most recent bot reply.
* Returns null if no such message exists (the user is "caught up").
*
* Used by restart-recovery to rewind the agent cursor before re-running an
* interrupted turn from the message that initiated it.
*/
export function getEarliestUnansweredHumanSeqFromDatabase(
export function getRecentChatMessagesBatchFromDatabase(
database: Database,
chatJid: string,
): number | null {
const lastBot = database
.prepare(
`SELECT MAX(seq) AS seq FROM messages
WHERE chat_jid = ? AND is_bot_message = 1`,
chatJids: string[],
limit: number = 8,
): Map<string, NewMessage[]> {
const out = new Map<string, NewMessage[]>();
if (chatJids.length === 0) return out;
const placeholders = chatJids.map(() => '?').join(',');
const sql = `
WITH ranked AS (
SELECT id, chat_jid, sender, sender_name, content, timestamp,
is_from_me, is_bot_message, message_source_kind,
ROW_NUMBER() OVER (PARTITION BY chat_jid ORDER BY timestamp DESC) AS rn
FROM messages
WHERE chat_jid IN (${placeholders})
AND content != '' AND content IS NOT NULL
)
.get(chatJid) as { seq: number | null } | undefined;
const lastBotSeq = lastBot?.seq ?? 0;
const row = database
.prepare(
`SELECT MIN(seq) AS seq FROM messages
WHERE chat_jid = ?
AND seq > ?
AND is_bot_message = 0
AND is_from_me = 0
AND content != '' AND content IS NOT NULL`,
)
.get(chatJid, lastBotSeq) as { seq: number | null } | undefined;
return row?.seq ?? null;
SELECT id, chat_jid, sender, sender_name, content, timestamp,
is_from_me, is_bot_message, message_source_kind
FROM ranked
WHERE rn <= ?
ORDER BY chat_jid, timestamp ASC
`;
const rows = database.prepare(sql).all(...chatJids, limit) as Array<
NewMessage & {
is_from_me?: boolean | number;
is_bot_message?: boolean | number;
}
>;
for (const row of rows) {
const normalized = normalizeMessageRow(row);
const existing = out.get(normalized.chat_jid);
if (existing) existing.push(normalized);
else out.set(normalized.chat_jid, [normalized]);
}
return out;
}
export function getLastHumanMessageTimestampFromDatabase(

View File

@@ -1,23 +0,0 @@
import type { Database } from 'bun:sqlite';
import { tableHasColumn } from './helpers.js';
import type { SchemaMigrationDefinition } from './types.js';
export const REVIEWER_FAILURE_COUNT_MIGRATION: SchemaMigrationDefinition = {
version: 15,
name: 'reviewer_failure_count',
apply(database: Database) {
if (!tableHasColumn(database, 'paired_tasks', 'reviewer_failure_count')) {
database.exec(`
ALTER TABLE paired_tasks
ADD COLUMN reviewer_failure_count INTEGER NOT NULL DEFAULT 0
`);
}
database.exec(`
UPDATE paired_tasks
SET reviewer_failure_count = 0
WHERE reviewer_failure_count IS NULL
`);
},
};

View File

@@ -0,0 +1,19 @@
import type { Database } from 'bun:sqlite';
import { tableHasColumn } from './helpers.js';
import type { SchemaMigrationDefinition } from './types.js';
export const TURN_PROGRESS_TEXT_MIGRATION: SchemaMigrationDefinition = {
version: 15,
name: 'turn_progress_text',
apply(database: Database) {
if (!tableHasColumn(database, 'paired_turns', 'progress_text')) {
database.exec(`ALTER TABLE paired_turns ADD COLUMN progress_text TEXT`);
}
if (!tableHasColumn(database, 'paired_turns', 'progress_updated_at')) {
database.exec(
`ALTER TABLE paired_turns ADD COLUMN progress_updated_at TEXT`,
);
}
},
};

View File

@@ -0,0 +1,29 @@
import type { Database } from 'bun:sqlite';
import type { SchemaMigrationDefinition } from './types.js';
export const ROOM_SKILL_OVERRIDES_MIGRATION: SchemaMigrationDefinition = {
version: 16,
name: 'room_skill_overrides',
apply(database: Database) {
database.exec(`
CREATE TABLE IF NOT EXISTS room_skill_overrides (
chat_jid TEXT NOT NULL,
agent_type TEXT NOT NULL,
skill_scope TEXT NOT NULL,
skill_name TEXT NOT NULL,
enabled INTEGER NOT NULL,
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (chat_jid, agent_type, skill_scope, skill_name),
FOREIGN KEY (chat_jid) REFERENCES room_settings(chat_jid) ON DELETE CASCADE,
CHECK (agent_type IN ('claude-code', 'codex')),
CHECK (enabled IN (0, 1)),
CHECK (length(skill_scope) > 0),
CHECK (length(skill_name) > 0)
);
CREATE INDEX IF NOT EXISTS idx_room_skill_overrides_room
ON room_skill_overrides(chat_jid, agent_type);
`);
},
};

View File

@@ -0,0 +1,13 @@
import type { SchemaMigrationDefinition } from './types.js';
import { tryExecMigration } from './helpers.js';
export const SCHEDULED_TASK_ROOM_ROLE_MIGRATION = {
version: 17,
name: 'scheduled_task_room_role',
apply(database) {
tryExecMigration(
database,
`ALTER TABLE scheduled_tasks ADD COLUMN room_role TEXT`,
);
},
} satisfies SchemaMigrationDefinition;

View File

@@ -0,0 +1,20 @@
import type { Database } from 'bun:sqlite';
import { tableHasColumn } from './helpers.js';
import type { SchemaMigrationDefinition } from './types.js';
export const PAIRED_TURN_OUTPUT_ATTACHMENTS_MIGRATION: SchemaMigrationDefinition =
{
version: 18,
name: 'paired_turn_output_attachments',
apply(database: Database) {
if (
!tableHasColumn(database, 'paired_turn_outputs', 'attachment_payload')
) {
database.exec(`
ALTER TABLE paired_turn_outputs
ADD COLUMN attachment_payload TEXT
`);
}
},
};

View File

@@ -14,7 +14,10 @@ import { OWNER_FAILURE_COUNT_MIGRATION } from './011_owner-failure-count.js';
import { PAIRED_VERDICT_AND_STEP_TELEMETRY_MIGRATION } from './012_paired-verdict-and-step-telemetry.js';
import { MESSAGE_SOURCE_KIND_MIGRATION } from './013_message-source-kind.js';
import { WORK_ITEM_ATTACHMENTS_MIGRATION } from './014_work-item-attachments.js';
import { REVIEWER_FAILURE_COUNT_MIGRATION } from './015_reviewer-failure-count.js';
import { TURN_PROGRESS_TEXT_MIGRATION } from './015_turn-progress-text.js';
import { ROOM_SKILL_OVERRIDES_MIGRATION } from './016_room-skill-overrides.js';
import { SCHEDULED_TASK_ROOM_ROLE_MIGRATION } from './017_scheduled-task-room-role.js';
import { PAIRED_TURN_OUTPUT_ATTACHMENTS_MIGRATION } from './018_paired-turn-output-attachments.js';
import type {
SchemaMigrationArgs,
SchemaMigrationDefinition,
@@ -37,7 +40,10 @@ const ORDERED_SCHEMA_MIGRATIONS: readonly SchemaMigrationDefinition[] = [
PAIRED_VERDICT_AND_STEP_TELEMETRY_MIGRATION,
MESSAGE_SOURCE_KIND_MIGRATION,
WORK_ITEM_ATTACHMENTS_MIGRATION,
REVIEWER_FAILURE_COUNT_MIGRATION,
TURN_PROGRESS_TEXT_MIGRATION,
ROOM_SKILL_OVERRIDES_MIGRATION,
SCHEDULED_TASK_ROOM_ROLE_MIGRATION,
PAIRED_TURN_OUTPUT_ATTACHMENTS_MIGRATION,
];
function ensureSchemaMigrationsTable(database: Database): void {

View File

@@ -20,7 +20,6 @@ import {
import {
AgentType,
PairedProject,
PairedRoomRole,
PairedTask,
PairedTaskStatus,
PairedTurnReservationIntentKind,
@@ -51,7 +50,6 @@ export type PairedTaskUpdates = Partial<
| 'review_requested_at'
| 'round_trip_count'
| 'owner_failure_count'
| 'reviewer_failure_count'
| 'owner_step_done_streak'
| 'finalize_step_done_count'
| 'task_done_then_user_reopen_count'
@@ -178,7 +176,6 @@ export function createPairedTaskInDatabase(
review_requested_at,
round_trip_count,
owner_failure_count,
reviewer_failure_count,
owner_step_done_streak,
finalize_step_done_count,
task_done_then_user_reopen_count,
@@ -190,7 +187,7 @@ export function createPairedTaskInDatabase(
created_at,
updated_at
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`,
)
.run(
@@ -208,7 +205,6 @@ export function createPairedTaskInDatabase(
task.review_requested_at,
task.round_trip_count,
task.owner_failure_count ?? 0,
task.reviewer_failure_count ?? 0,
task.owner_step_done_streak ?? 0,
task.finalize_step_done_count ?? 0,
task.task_done_then_user_reopen_count ?? 0,
@@ -232,21 +228,27 @@ export function getPairedTaskByIdFromDatabase(
return row ? hydratePairedTaskRow(database, row) : undefined;
}
const latestPairedTaskStmtCache = new WeakMap<
Database,
ReturnType<Database['prepare']>
>();
export function getLatestPairedTaskForChatFromDatabase(
database: Database,
chatJid: string,
): PairedTask | undefined {
const row = database
.prepare(
`
SELECT *
FROM paired_tasks
WHERE chat_jid = ?
ORDER BY updated_at DESC
LIMIT 1
`,
)
.get(chatJid) as StoredPairedTaskRow | undefined;
let stmt = latestPairedTaskStmtCache.get(database);
if (!stmt) {
stmt = database.prepare(`
SELECT *
FROM paired_tasks
WHERE chat_jid = ?
ORDER BY updated_at DESC
LIMIT 1
`);
latestPairedTaskStmtCache.set(database, stmt);
}
const row = stmt.get(chatJid) as StoredPairedTaskRow | undefined;
return row ? hydratePairedTaskRow(database, row) : undefined;
}
@@ -269,6 +271,22 @@ export function getLatestOpenPairedTaskForChatFromDatabase(
return row ? hydratePairedTaskRow(database, row) : undefined;
}
export function getAllOpenPairedTasksFromDatabase(
database: Database,
): PairedTask[] {
const rows = database
.prepare(
`
SELECT *
FROM paired_tasks
WHERE status NOT IN ('completed')
ORDER BY updated_at DESC, created_at DESC
`,
)
.all() as StoredPairedTaskRow[];
return rows.map((row) => hydratePairedTaskRow(database, row));
}
export function getLatestPreviousPairedTaskForChatFromDatabase(
database: Database,
chatJid: string,
@@ -321,10 +339,6 @@ export function updatePairedTaskInDatabase(
fields.push('owner_failure_count = ?');
values.push(updates.owner_failure_count);
}
if (updates.reviewer_failure_count !== undefined) {
fields.push('reviewer_failure_count = ?');
values.push(updates.reviewer_failure_count);
}
if (updates.owner_step_done_streak !== undefined) {
fields.push('owner_step_done_streak = ?');
values.push(updates.owner_step_done_streak);
@@ -404,10 +418,6 @@ export function updatePairedTaskIfUnchangedInDatabase(
fields.push('owner_failure_count = ?');
values.push(updates.owner_failure_count);
}
if (updates.reviewer_failure_count !== undefined) {
fields.push('reviewer_failure_count = ?');
values.push(updates.reviewer_failure_count);
}
if (updates.owner_step_done_streak !== undefined) {
fields.push('owner_step_done_streak = ?');
values.push(updates.owner_step_done_streak);

View File

@@ -1,6 +1,7 @@
import { Database } from 'bun:sqlite';
import { normalizeServiceId } from '../config.js';
import { CODEX_BAD_REQUEST_DETAIL_JSON } from '../codex-bad-request-signal.js';
import type { PairedTurnIdentity } from '../paired-turn-identity.js';
import { inferAgentTypeFromServiceShadow } from '../role-service-shadow.js';
import type { AgentType, PairedRoomRole } from '../types.js';
@@ -33,6 +34,25 @@ export interface PairedTurnAttemptRecord {
last_error: string | null;
}
export interface InterruptedPairedTurnAttemptRecoveryCandidate {
chat_jid: string;
group_folder: string;
task_id: string;
task_status: string;
turn_id: string;
attempt_id: string;
attempt_no: number;
role: PairedRoomRole;
intent_kind: PairedTurnIdentity['intentKind'];
}
export interface OwnerCodexBadRequestFailureSummary {
taskId: string;
failures: number;
firstFailureAt: string;
latestFailureAt: string;
}
function resolveExecutorMetadata(args: {
executorServiceId?: string | null;
executorAgentType?: AgentType | null;
@@ -423,6 +443,154 @@ export function getPairedTurnAttemptIdFromDatabase(
return row?.attempt_id ?? null;
}
export function recoverInterruptedPairedTurnAttemptsForServiceInDatabase(
database: Database,
args: {
serviceIds: string[];
now?: string;
error?: string;
},
): InterruptedPairedTurnAttemptRecoveryCandidate[] {
const serviceIds = Array.from(
new Set(args.serviceIds.map((serviceId) => normalizeServiceId(serviceId))),
).filter((serviceId) => serviceId.length > 0);
if (serviceIds.length === 0) {
return [];
}
const servicePlaceholders = serviceIds.map(() => '?').join(', ');
const now = args.now ?? new Date().toISOString();
const error =
args.error ?? 'Interrupted by service restart before completion.';
return database.transaction(() => {
database
.prepare(
`
UPDATE paired_turn_attempts
SET state = 'failed',
active_run_id = NULL,
updated_at = ?,
completed_at = ?,
last_error = COALESCE(last_error, ?)
WHERE state = 'running'
AND executor_service_id IN (${servicePlaceholders})
AND task_id IN (
SELECT id FROM paired_tasks WHERE status = 'completed'
)
`,
)
.run(now, now, error, ...serviceIds);
const rows = database
.prepare(
`
SELECT
tasks.chat_jid,
tasks.group_folder,
attempts.task_id,
tasks.status AS task_status,
attempts.turn_id,
attempts.attempt_id,
attempts.attempt_no,
attempts.role,
attempts.intent_kind
FROM paired_turn_attempts attempts
JOIN paired_tasks tasks
ON tasks.id = attempts.task_id
WHERE attempts.state = 'running'
AND attempts.executor_service_id IN (${servicePlaceholders})
AND tasks.status != 'completed'
AND NOT EXISTS (
SELECT 1
FROM paired_task_execution_leases leases
WHERE leases.task_id = attempts.task_id
AND leases.turn_id = attempts.turn_id
AND leases.claimed_service_id NOT IN (${servicePlaceholders})
)
ORDER BY attempts.updated_at ASC, attempts.attempt_no ASC
`,
)
.all(
...serviceIds,
...serviceIds,
) as InterruptedPairedTurnAttemptRecoveryCandidate[];
if (rows.length === 0) {
return rows;
}
const update = database.prepare(
`
UPDATE paired_turn_attempts
SET state = 'failed',
active_run_id = NULL,
updated_at = ?,
completed_at = ?,
last_error = ?
WHERE attempt_id = ?
AND state = 'running'
AND executor_service_id IN (${servicePlaceholders})
`,
);
for (const row of rows) {
update.run(now, now, error, row.attempt_id, ...serviceIds);
}
return rows;
})();
}
export function getOwnerCodexBadRequestFailureSummaryForTaskFromDatabase(
database: Database,
args: {
taskId: string;
threshold: number;
},
): OwnerCodexBadRequestFailureSummary | null {
const threshold = Math.max(1, Math.floor(args.threshold));
const attempts = database
.prepare(
`
SELECT state, last_error, created_at
FROM paired_turn_attempts
WHERE task_id = ?
AND role = 'owner'
AND executor_agent_type = 'codex'
ORDER BY created_at DESC, attempt_no DESC
`,
)
.all(args.taskId) as Array<{
state: PairedTurnAttemptState;
last_error: string | null;
created_at: string;
}>;
const consecutiveFailures = [];
for (const attempt of attempts) {
if (
attempt.state === 'failed' &&
attempt.last_error?.trim() === CODEX_BAD_REQUEST_DETAIL_JSON
) {
consecutiveFailures.push(attempt);
continue;
}
break;
}
if (consecutiveFailures.length < threshold) {
return null;
}
const firstFailureAt =
consecutiveFailures[consecutiveFailures.length - 1].created_at;
const latestFailureAt = consecutiveFailures[0].created_at;
return {
taskId: args.taskId,
failures: consecutiveFailures.length,
firstFailureAt,
latestFailureAt,
};
}
export function setPairedTurnAttemptContinuationHandoffIdInDatabase(
database: Database,
args: {

View File

@@ -1,21 +1,38 @@
import { Database } from 'bun:sqlite';
import { logger } from '../logger.js';
import { parseVisibleVerdict } from '../paired-verdict.js';
import {
parseReviewerVisibleVerdict,
parseVisibleVerdict,
} from '../paired-verdict.js';
import { PairedRoomRole, PairedTurnOutput } from '../types.js';
OutboundAttachment,
PairedRoomRole,
PairedTurnOutput,
} from '../types.js';
import {
parseAttachmentPayload,
serializeAttachmentPayload,
} from './work-items.js';
const MAX_TURN_OUTPUT_CHARS = 50_000;
function storedOutputText(outputText: string): string {
if (outputText.length <= MAX_TURN_OUTPUT_CHARS) {
return outputText;
}
const notice = `\n\n[Output truncated: ${outputText.length} > ${MAX_TURN_OUTPUT_CHARS} chars]`;
return `${outputText.slice(0, MAX_TURN_OUTPUT_CHARS - notice.length)}${notice}`;
}
export function insertPairedTurnOutputInDatabase(
database: Database,
taskId: string,
turnNumber: number,
role: PairedRoomRole,
outputText: string,
createdAt?: string,
options: {
createdAt?: string;
attachments?: OutboundAttachment[];
} = {},
): void {
if (outputText.length > MAX_TURN_OUTPUT_CHARS) {
logger.warn(
@@ -33,21 +50,36 @@ export function insertPairedTurnOutputInDatabase(
database
.prepare(
`INSERT OR REPLACE INTO paired_turn_outputs
(task_id, turn_number, role, output_text, verdict, created_at)
VALUES (?, ?, ?, ?, ?, ?)`,
(task_id, turn_number, role, output_text, attachment_payload, verdict, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?)`,
)
.run(
taskId,
turnNumber,
role,
outputText.slice(0, MAX_TURN_OUTPUT_CHARS),
role === 'reviewer'
? parseReviewerVisibleVerdict(outputText)
: parseVisibleVerdict(outputText),
createdAt ?? new Date().toISOString(),
storedOutputText(outputText),
serializeAttachmentPayload(options.attachments),
parseVisibleVerdict(outputText),
options.createdAt ?? new Date().toISOString(),
);
}
type StoredPairedTurnOutputRow = PairedTurnOutput & {
attachment_payload?: string | null;
};
function hydratePairedTurnOutputRow(
row: StoredPairedTurnOutputRow,
): PairedTurnOutput {
return {
...row,
attachments: parseAttachmentPayload(row.attachment_payload, {
table: 'paired_turn_outputs',
rowId: row.id,
}),
};
}
export function getPairedTurnOutputsFromDatabase(
database: Database,
taskId: string,
@@ -58,7 +90,34 @@ export function getPairedTurnOutputsFromDatabase(
WHERE task_id = ?
ORDER BY turn_number ASC`,
)
.all(taskId) as PairedTurnOutput[];
.all(taskId)
.map((row) => hydratePairedTurnOutputRow(row as StoredPairedTurnOutputRow));
}
const recentOutputsForChatStmtCache = new WeakMap<
Database,
ReturnType<Database['prepare']>
>();
export function getRecentPairedTurnOutputsForChatFromDatabase(
database: Database,
chatJid: string,
limit: number = 8,
): PairedTurnOutput[] {
let stmt = recentOutputsForChatStmtCache.get(database);
if (!stmt) {
stmt = database.prepare(`
SELECT o.*
FROM paired_turn_outputs o
INNER JOIN paired_tasks t ON o.task_id = t.id
WHERE t.chat_jid = ?
ORDER BY o.created_at DESC
LIMIT ?
`);
recentOutputsForChatStmtCache.set(database, stmt);
}
const rows = stmt.all(chatJid, limit) as StoredPairedTurnOutputRow[];
return rows.reverse().map(hydratePairedTurnOutputRow);
}
export function getLatestTurnNumberFromDatabase(

View File

@@ -1,7 +1,7 @@
import { Database } from 'bun:sqlite';
import { buildPairedTurnAttemptId } from './paired-turn-attempts.js';
import { tableHasColumn, tryExecMigration } from './migrations/helpers.js';
import { tableHasColumn } from './migrations/helpers.js';
// Paired-turn provenance rebuild helpers extracted from the legacy schema
// bundle. These remain runtime helpers because v10 replays them during

View File

@@ -33,6 +33,8 @@ export interface PairedTurnRecord {
updated_at: string;
completed_at: string | null;
last_error: string | null;
progress_text?: string | null;
progress_updated_at?: string | null;
}
interface StoredPairedTurnRow {
@@ -43,6 +45,8 @@ interface StoredPairedTurnRow {
intent_kind: PairedTurnIdentity['intentKind'];
created_at: string;
updated_at: string;
progress_text?: string | null;
progress_updated_at?: string | null;
}
function hydratePairedTurnRecord(
@@ -493,6 +497,67 @@ export function getPairedTurnByIdFromDatabase(
);
}
const updateProgressTextStmtCache = new WeakMap<
Database,
ReturnType<Database['prepare']>
>();
export function updatePairedTurnProgressTextFromDatabase(
database: Database,
turnId: string,
progressText: string | null,
): void {
const now = new Date().toISOString();
let stmt = updateProgressTextStmtCache.get(database);
if (!stmt) {
stmt = database.prepare(`
UPDATE paired_turns
SET progress_text = ?,
progress_updated_at = ?,
updated_at = ?
WHERE turn_id = ?
`);
updateProgressTextStmtCache.set(database, stmt);
}
stmt.run(progressText, now, now, turnId);
}
const latestPairedTurnStmtCache = new WeakMap<
Database,
ReturnType<Database['prepare']>
>();
export function getLatestPairedTurnForTaskFromDatabase(
database: Database,
taskId: string,
): PairedTurnRecord | null {
let stmt = latestPairedTurnStmtCache.get(database);
if (!stmt) {
stmt = database.prepare(`
SELECT *
FROM paired_turns
WHERE task_id = ?
ORDER BY CASE
WHEN progress_text IS NOT NULL
AND trim(progress_text) <> ''
AND progress_updated_at IS NOT NULL
AND progress_updated_at > updated_at
THEN progress_updated_at
ELSE updated_at
END DESC,
turn_id DESC
LIMIT 1
`);
latestPairedTurnStmtCache.set(database, stmt);
}
const row = stmt.get(taskId) as StoredPairedTurnRow | undefined;
if (!row) return null;
return hydratePairedTurnRecord(
row,
getCurrentPairedTurnAttemptForTurnFromDatabase(database, row.turn_id),
);
}
export function getPairedTurnsForTaskFromDatabase(
database: Database,
taskId: string,

View File

@@ -115,7 +115,7 @@ function getStoredRoomRoleOverrideRows(
agent_config_json: string | null;
created_at: string;
updated_at: string;
}> = [];
}>;
try {
rows = database
.prepare(

View File

@@ -28,6 +28,24 @@ interface StoredRoomModeRow {
source: RoomModeSource;
}
export interface StoredRoomSkillOverride {
chatJid: string;
agentType: AgentType;
skillScope: string;
skillName: string;
enabled: boolean;
createdAt: string;
updatedAt: string;
}
export interface StoredRoomSkillOverrideInput {
chatJid: string;
agentType: AgentType;
skillScope: string;
skillName: string;
enabled: boolean;
}
export interface AssignRoomInput {
name: string;
roomMode?: RoomMode;
@@ -35,6 +53,8 @@ export interface AssignRoomInput {
reviewerAgentType?: AgentType;
arbiterAgentType?: AgentType | null;
folder?: string;
trigger?: string;
requiresTrigger?: boolean;
isMain?: boolean;
workDir?: string;
addedAt?: string;
@@ -186,8 +206,9 @@ export function assignRoomInDatabase(
const snapshot: RoomRegistrationSnapshot = {
name: input.name,
folder,
triggerPattern: existing?.trigger ?? '',
requiresTrigger: existing?.requiresTrigger ?? false,
triggerPattern: input.trigger ?? existing?.trigger ?? '',
requiresTrigger:
input.requiresTrigger ?? existing?.requiresTrigger ?? false,
isMain: input.isMain ?? existing?.isMain ?? false,
ownerAgentType,
workDir: input.workDir ?? existing?.workDir ?? null,
@@ -307,6 +328,112 @@ export function getStoredRoomRoleAgentPlanFromDatabase(
return stored ? resolveStoredRoomRoleAgentPlan(database, stored) : undefined;
}
export function getStoredRoomSkillOverridesFromDatabase(
database: Database,
chatJid?: string,
): StoredRoomSkillOverride[] {
const params: string[] = [];
const where = chatJid ? 'WHERE chat_jid = ?' : '';
if (chatJid) params.push(chatJid);
let rows: Array<{
chat_jid: string;
agent_type: string;
skill_scope: string;
skill_name: string;
enabled: number;
created_at: string;
updated_at: string;
}>;
try {
rows = database
.prepare(
`SELECT chat_jid, agent_type, skill_scope, skill_name, enabled,
created_at, updated_at
FROM room_skill_overrides
${where}
ORDER BY chat_jid, agent_type, skill_scope, skill_name`,
)
.all(...params) as Array<{
chat_jid: string;
agent_type: string;
skill_scope: string;
skill_name: string;
enabled: number;
created_at: string;
updated_at: string;
}>;
} catch {
return [];
}
return rows
.map((row) => {
const agentType = normalizeStoredAgentType(row.agent_type);
if (!agentType) return null;
return {
chatJid: row.chat_jid,
agentType,
skillScope: row.skill_scope,
skillName: row.skill_name,
enabled: row.enabled === 1,
createdAt: row.created_at,
updatedAt: row.updated_at,
};
})
.filter((row): row is StoredRoomSkillOverride => Boolean(row));
}
export function upsertStoredRoomSkillOverrideInDatabase(
database: Database,
input: StoredRoomSkillOverrideInput,
): void {
const agentType = normalizeStoredAgentType(input.agentType);
if (!agentType) {
throw new Error(`Unsupported agent type: ${input.agentType}`);
}
const now = new Date().toISOString();
database
.prepare(
`INSERT INTO room_skill_overrides (
chat_jid, agent_type, skill_scope, skill_name, enabled,
created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(chat_jid, agent_type, skill_scope, skill_name)
DO UPDATE SET
enabled = excluded.enabled,
updated_at = excluded.updated_at`,
)
.run(
input.chatJid,
agentType,
input.skillScope,
input.skillName,
input.enabled ? 1 : 0,
now,
now,
);
}
export function deleteStoredRoomSkillOverrideFromDatabase(
database: Database,
input: Omit<StoredRoomSkillOverrideInput, 'enabled'>,
): void {
const agentType = normalizeStoredAgentType(input.agentType);
if (!agentType) {
throw new Error(`Unsupported agent type: ${input.agentType}`);
}
database
.prepare(
`DELETE FROM room_skill_overrides
WHERE chat_jid = ?
AND agent_type = ?
AND skill_scope = ?
AND skill_name = ?`,
)
.run(input.chatJid, agentType, input.skillScope, input.skillName);
}
function getStoredRoomModeRowFromDatabase(
database: Database,
chatJid: string,

View File

@@ -17,7 +17,6 @@ import {
import {
type ChatInfo,
getAllChatsFromDatabase,
getEarliestUnansweredHumanSeqFromDatabase,
getLastHumanMessageContentFromDatabase,
getLastHumanMessageSenderFromDatabase,
getLastHumanMessageTimestampFromDatabase,
@@ -27,6 +26,8 @@ import {
getNewMessagesBySeqFromDatabase,
getNewMessagesFromDatabase,
getRecentChatMessagesFromDatabase,
getRecentChatMessagesBatchFromDatabase,
hasMessageInDatabase,
hasRecentRestartAnnouncementInDatabase,
storeChatMetadataInDatabase,
storeMessageInDatabase,
@@ -37,7 +38,7 @@ import {
createProducedWorkItemInDatabase,
getOpenWorkItemForChatFromDatabase,
getOpenWorkItemFromDatabase,
getRecentDeliveredOwnerWorkItemsForChatFromDatabase,
getRecentDeliveredWorkItemsForChatFromDatabase,
markWorkItemDeliveredInDatabase,
markWorkItemDeliveryRetryInDatabase,
} from './work-items.js';
@@ -143,6 +144,10 @@ export function storeMessage(msg: NewMessage): void {
storeMessageInDatabase(requireDatabase(), msg);
}
export function hasMessage(chatJid: string, id: string): boolean {
return hasMessageInDatabase(requireDatabase(), chatJid, id);
}
export function getNewMessages(
jids: string[],
lastTimestamp: string,
@@ -221,12 +226,19 @@ export function getRecentChatMessages(
return getRecentChatMessagesFromDatabase(requireDatabase(), chatJid, limit);
}
export function getLastHumanMessageTimestamp(chatJid: string): string | null {
return getLastHumanMessageTimestampFromDatabase(requireDatabase(), chatJid);
export function getRecentChatMessagesBatch(
chatJids: string[],
limit: number = 8,
): Map<string, NewMessage[]> {
return getRecentChatMessagesBatchFromDatabase(
requireDatabase(),
chatJids,
limit,
);
}
export function getEarliestUnansweredHumanSeq(chatJid: string): number | null {
return getEarliestUnansweredHumanSeqFromDatabase(requireDatabase(), chatJid);
export function getLastHumanMessageTimestamp(chatJid: string): string | null {
return getLastHumanMessageTimestampFromDatabase(requireDatabase(), chatJid);
}
export function getLastHumanMessageSender(chatJid: string): string | null {
@@ -272,11 +284,11 @@ export function getOpenWorkItemForChat(
);
}
export function getRecentDeliveredOwnerWorkItemsForChat(
export function getRecentDeliveredWorkItemsForChat(
chatJid: string,
limit: number,
limit: number = 8,
): WorkItem[] {
return getRecentDeliveredOwnerWorkItemsForChatFromDatabase(
return getRecentDeliveredWorkItemsForChatFromDatabase(
requireDatabase(),
chatJid,
limit,

View File

@@ -24,6 +24,7 @@ import {
clearPairedTurnReservationsInDatabase,
type PairedTaskUpdates,
createPairedTaskInDatabase,
getAllOpenPairedTasksFromDatabase,
getLastBotFinalMessageFromDatabase,
getLatestOpenPairedTaskForChatFromDatabase,
getLatestPreviousPairedTaskForChatFromDatabase,
@@ -42,12 +43,17 @@ import {
} from './paired-state.js';
import {
clearPairedTurnAttemptsInDatabase,
getOwnerCodexBadRequestFailureSummaryForTaskFromDatabase,
getPairedTurnAttemptsForTurnFromDatabase,
recoverInterruptedPairedTurnAttemptsForServiceInDatabase,
type InterruptedPairedTurnAttemptRecoveryCandidate,
type OwnerCodexBadRequestFailureSummary,
type PairedTurnAttemptRecord,
} from './paired-turn-attempts.js';
import {
getLatestTurnNumberFromDatabase,
getPairedTurnOutputsFromDatabase,
getRecentPairedTurnOutputsForChatFromDatabase,
insertPairedTurnOutputInDatabase,
} from './paired-turn-outputs.js';
import {
@@ -57,6 +63,8 @@ import {
failPairedTurnInDatabase,
getPairedTurnByIdFromDatabase,
getPairedTurnsForTaskFromDatabase,
getLatestPairedTurnForTaskFromDatabase,
updatePairedTurnProgressTextFromDatabase,
markPairedTurnRunningInDatabase,
type PairedTurnRecord,
} from './paired-turns.js';
@@ -113,6 +121,10 @@ export function getLatestPreviousPairedTaskForChat(
);
}
export function getAllOpenPairedTasks(): PairedTask[] {
return getAllOpenPairedTasksFromDatabase(requireDatabase());
}
export function updatePairedTask(id: string, updates: PairedTaskUpdates): void {
updatePairedTaskInDatabase(requireDatabase(), id, updates);
}
@@ -206,12 +218,50 @@ export function getPairedTurnsForTask(taskId: string): PairedTurnRecord[] {
return getPairedTurnsForTaskFromDatabase(requireDatabase(), taskId);
}
export function getLatestPairedTurnForTask(
taskId: string,
): PairedTurnRecord | null {
return getLatestPairedTurnForTaskFromDatabase(requireDatabase(), taskId);
}
export function updatePairedTurnProgressText(
turnId: string,
progressText: string | null,
): void {
updatePairedTurnProgressTextFromDatabase(
requireDatabase(),
turnId,
progressText,
);
}
export function getPairedTurnAttempts(
turnId: string,
): PairedTurnAttemptRecord[] {
return getPairedTurnAttemptsForTurnFromDatabase(requireDatabase(), turnId);
}
export function recoverInterruptedPairedTurnAttemptsForService(args: {
serviceIds: string[];
now?: string;
error?: string;
}): InterruptedPairedTurnAttemptRecoveryCandidate[] {
return recoverInterruptedPairedTurnAttemptsForServiceInDatabase(
requireDatabase(),
args,
);
}
export function getOwnerCodexBadRequestFailureSummaryForTask(args: {
taskId: string;
threshold: number;
}): OwnerCodexBadRequestFailureSummary | null {
return getOwnerCodexBadRequestFailureSummaryForTaskFromDatabase(
requireDatabase(),
args,
);
}
export function upsertPairedWorkspace(workspace: PairedWorkspace): void {
upsertPairedWorkspaceInDatabase(requireDatabase(), workspace);
}
@@ -312,15 +362,24 @@ export function insertPairedTurnOutput(
turnNumber: number,
role: PairedRoomRole,
outputText: string,
createdAt?: string,
createdAtOrOptions?:
| string
| {
createdAt?: string;
attachments?: import('../types.js').OutboundAttachment[];
},
): void {
const options =
typeof createdAtOrOptions === 'string'
? { createdAt: createdAtOrOptions }
: (createdAtOrOptions ?? {});
insertPairedTurnOutputInDatabase(
requireDatabase(),
taskId,
turnNumber,
role,
outputText,
createdAt,
options,
);
}
@@ -328,6 +387,17 @@ export function getPairedTurnOutputs(taskId: string): PairedTurnOutput[] {
return getPairedTurnOutputsFromDatabase(requireDatabase(), taskId);
}
export function getRecentPairedTurnOutputsForChat(
chatJid: string,
limit: number = 8,
): PairedTurnOutput[] {
return getRecentPairedTurnOutputsForChatFromDatabase(
requireDatabase(),
chatJid,
limit,
);
}
export function getLatestTurnNumber(taskId: string): number {
return getLatestTurnNumberFromDatabase(requireDatabase(), taskId);
}

View File

@@ -28,6 +28,7 @@ import {
assignRoomInDatabase,
clearExplicitRoomModeInDatabase,
deleteStoredRoomSettingsForTestsInDatabase,
deleteStoredRoomSkillOverrideFromDatabase,
getAllRoomBindingsFromDatabase,
getEffectiveRoomModeFromDatabase,
getEffectiveRuntimeRoomModeFromDatabase,
@@ -36,10 +37,14 @@ import {
getRegisteredGroupFromDatabase,
getStoredRoomRoleAgentPlanFromDatabase,
getStoredRoomSettingsFromDatabase,
getStoredRoomSkillOverridesFromDatabase,
setExplicitRoomModeInDatabase,
setRegisteredGroupForTestsInDatabase,
setStoredRoomOwnerAgentTypeForTestsInDatabase,
type StoredRoomSkillOverride,
type StoredRoomSkillOverrideInput,
updateRegisteredGroupNameInDatabase,
upsertStoredRoomSkillOverrideInDatabase,
} from './rooms.js';
import { type StoredRoomSettings } from './room-registration.js';
import {
@@ -363,6 +368,26 @@ export function getStoredRoomRoleAgentPlan(
return getStoredRoomRoleAgentPlanFromDatabase(db, chatJid);
}
export function getStoredRoomSkillOverrides(
chatJid?: string,
): StoredRoomSkillOverride[] {
const db = getDatabaseIfInitialized();
if (!db) return [];
return getStoredRoomSkillOverridesFromDatabase(db, chatJid);
}
export function upsertStoredRoomSkillOverride(
input: StoredRoomSkillOverrideInput,
): void {
upsertStoredRoomSkillOverrideInDatabase(requireDatabase(), input);
}
export function deleteStoredRoomSkillOverride(
input: Omit<StoredRoomSkillOverrideInput, 'enabled'>,
): void {
deleteStoredRoomSkillOverrideFromDatabase(requireDatabase(), input);
}
export function getExplicitRoomMode(chatJid: string): RoomMode | undefined {
return getExplicitRoomModeFromDatabase(requireDatabase(), chatJid);
}

View File

@@ -211,6 +211,7 @@ function parseLastAgentSeqState(
`Invalid last_agent_seq JSON for ${serviceId}: ${
err instanceof Error ? err.message : String(err)
}`,
{ cause: err },
);
}

View File

@@ -141,8 +141,11 @@ export function deleteAllSessionsForGroupFromDatabase(
groupFolder: string,
): void {
database
.prepare('DELETE FROM sessions WHERE group_folder = ?')
.run(groupFolder);
.prepare(
`DELETE FROM sessions
WHERE group_folder IN (?, ?, ?)`,
)
.run(groupFolder, `${groupFolder}:reviewer`, `${groupFolder}:arbiter`);
}
export function getAllSessionsForAgentTypeFromDatabase(

View File

@@ -1,12 +1,22 @@
import { Database } from 'bun:sqlite';
import {
DEFAULT_TASK_CONTEXT_MODE,
WATCH_CI_PROMPT_PREFIX,
} from 'ejclaw-runners-shared';
import { AgentType, ScheduledTask, TaskRunLog } from '../types.js';
import {
AgentType,
PairedRoomRole,
ScheduledTask,
TaskRunLog,
} from '../types.js';
export type CreateScheduledTaskInput = Omit<
ScheduledTask,
| 'last_run'
| 'last_result'
| 'agent_type'
| 'room_role'
| 'ci_provider'
| 'ci_metadata'
| 'max_duration_ms'
@@ -14,6 +24,7 @@ export type CreateScheduledTaskInput = Omit<
| 'status_started_at'
> & {
agent_type?: AgentType | null;
room_role?: PairedRoomRole | null;
ci_provider?: ScheduledTask['ci_provider'];
ci_metadata?: string | null;
max_duration_ms?: number | null;
@@ -45,8 +56,8 @@ export function createTaskInDatabase(
database
.prepare(
`
INSERT INTO scheduled_tasks (id, group_folder, chat_jid, agent_type, ci_provider, ci_metadata, max_duration_ms, status_message_id, status_started_at, prompt, schedule_type, schedule_value, context_mode, next_run, status, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
INSERT INTO scheduled_tasks (id, group_folder, chat_jid, agent_type, room_role, ci_provider, ci_metadata, max_duration_ms, status_message_id, status_started_at, prompt, schedule_type, schedule_value, context_mode, next_run, status, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`,
)
.run(
@@ -54,6 +65,7 @@ export function createTaskInDatabase(
task.group_folder,
task.chat_jid,
task.agent_type || 'claude-code',
task.room_role ?? null,
task.ci_provider ?? null,
task.ci_metadata ?? null,
task.max_duration_ms ?? null,
@@ -62,7 +74,7 @@ export function createTaskInDatabase(
task.prompt,
task.schedule_type,
task.schedule_value,
task.context_mode || 'isolated',
task.context_mode || DEFAULT_TASK_CONTEXT_MODE,
task.next_run,
task.status,
task.created_at,
@@ -73,9 +85,10 @@ export function getTaskByIdFromDatabase(
database: Database,
id: string,
): ScheduledTask | undefined {
return database
const row = database
.prepare('SELECT * FROM scheduled_tasks WHERE id = ?')
.get(id) as ScheduledTask | undefined;
.get(id) as ScheduledTask | null | undefined;
return row ?? undefined;
}
export function findDuplicateCiWatcherInDatabase(
@@ -208,10 +221,10 @@ export function hasActiveCiWatcherForChatInDatabase(
const row = database
.prepare(
`SELECT 1 FROM scheduled_tasks
WHERE chat_jid = ? AND status = 'active' AND prompt LIKE '[BACKGROUND CI WATCH]%'
WHERE chat_jid = ? AND status = 'active' AND prompt LIKE ?
LIMIT 1`,
)
.get(chatJid);
.get(chatJid, `${WATCH_CI_PROMPT_PREFIX}%`);
return !!row;
}

View File

@@ -1,6 +1,7 @@
import { Database } from 'bun:sqlite';
import { SERVICE_SESSION_SCOPE, normalizeServiceId } from '../config.js';
import { logger } from '../logger.js';
import {
inferAgentTypeFromServiceShadow,
inferRoleFromServiceShadow,
@@ -8,6 +9,8 @@ import {
} from '../role-service-shadow.js';
import { AgentType, OutboundAttachment, PairedRoomRole } from '../types.js';
const SUPERSEDED_WORK_ITEM_ERROR = 'superseded_by_newer_canonical_output';
export interface WorkItem {
id: number;
group_folder: string;
@@ -99,36 +102,107 @@ function hydrateWorkItemRow(row: StoredWorkItemRow): WorkItem {
...row,
agent_type: agentType,
service_id: readStoredWorkItemServiceId(row),
attachments: parseAttachmentPayload(row.attachment_payload),
attachments: parseAttachmentPayload(row.attachment_payload, {
table: 'work_items',
rowId: row.id,
}),
};
}
function parseAttachmentPayload(
export interface AttachmentPayloadContext {
table?: string;
rowId?: string | number;
}
function payloadPreview(payload: string): string {
return payload.length > 200 ? `${payload.slice(0, 200)}...` : payload;
}
function isAttachmentPayloadEntry(
item: unknown,
): item is { path: string; name?: unknown; mime?: unknown } {
return (
item !== null &&
typeof item === 'object' &&
!Array.isArray(item) &&
typeof (item as { path?: unknown }).path === 'string'
);
}
function logMalformedAttachmentPayload(args: {
payload: string;
context?: AttachmentPayloadContext;
reason: string;
err?: unknown;
}): void {
logger.warn(
{
...args.context,
reason: args.reason,
payloadLength: args.payload.length,
payloadPreview: payloadPreview(args.payload),
...(args.err ? { err: args.err } : {}),
},
'Ignored malformed attachment payload',
);
}
export function parseAttachmentPayload(
payload: string | null | undefined,
context?: AttachmentPayloadContext,
): OutboundAttachment[] {
if (!payload) return [];
try {
const parsed = JSON.parse(payload) as unknown;
if (!Array.isArray(parsed)) return [];
return parsed
.filter(
(item): item is OutboundAttachment =>
item !== null &&
typeof item === 'object' &&
!Array.isArray(item) &&
typeof (item as { path?: unknown }).path === 'string',
)
.map((item) => ({
path: item.path,
...(typeof item.name === 'string' ? { name: item.name } : {}),
...(typeof item.mime === 'string' ? { mime: item.mime } : {}),
}));
} catch {
if (!Array.isArray(parsed)) {
logMalformedAttachmentPayload({
payload,
context,
reason: 'not_array',
});
return [];
}
const attachments: OutboundAttachment[] = [];
let invalidEntryCount = 0;
for (const item of parsed) {
if (isAttachmentPayloadEntry(item)) {
attachments.push({
path: item.path,
...(typeof item.name === 'string' ? { name: item.name } : {}),
...(typeof item.mime === 'string' ? { mime: item.mime } : {}),
});
} else {
invalidEntryCount += 1;
}
}
if (invalidEntryCount > 0) {
logger.warn(
{
...context,
invalidEntryCount,
validEntryCount: attachments.length,
payloadLength: payload.length,
payloadPreview: payloadPreview(payload),
},
'Ignored invalid attachment payload entries',
);
}
return attachments;
} catch (err) {
logMalformedAttachmentPayload({
payload,
context,
reason: 'invalid_json',
err,
});
return [];
}
}
function serializeAttachmentPayload(
export function serializeAttachmentPayload(
attachments: OutboundAttachment[] | undefined,
): string | null {
if (!attachments?.length) return null;
@@ -267,6 +341,57 @@ export function getOpenWorkItemForChatFromDatabase(
return row ? hydrateWorkItemRow(row) : undefined;
}
export function getRecentDeliveredWorkItemsForChatFromDatabase(
database: Database,
chatJid: string,
limit: number = 8,
): WorkItem[] {
const rows = database
.prepare(
`SELECT *
FROM work_items
WHERE chat_jid = ?
AND status = 'delivered'
AND (last_error IS NULL OR last_error <> ?)
ORDER BY COALESCE(delivered_at, updated_at, created_at) DESC, id DESC
LIMIT ?`,
)
.all(chatJid, SUPERSEDED_WORK_ITEM_ERROR, limit) as StoredWorkItemRow[];
return rows.reverse().map(hydrateWorkItemRow);
}
function supersedeOpenWorkItemsForCanonicalKey(
database: Database,
input: CreateProducedWorkItemInput,
agentType: AgentType,
serviceId: string,
now: string,
): void {
database
.prepare(
`UPDATE work_items
SET status = 'delivered',
delivered_at = ?,
delivery_message_id = NULL,
last_error = ?,
updated_at = ?
WHERE chat_jid = ?
AND agent_type = ?
AND IFNULL(service_id, '') = IFNULL(?, '')
AND IFNULL(delivery_role, '') = IFNULL(?, '')
AND status IN ('produced', 'delivery_retry')`,
)
.run(
now,
SUPERSEDED_WORK_ITEM_ERROR,
now,
input.chat_jid,
agentType,
serviceId,
input.delivery_role ?? null,
);
}
export function createProducedWorkItemInDatabase(
database: Database,
input: CreateProducedWorkItemInput,
@@ -278,46 +403,57 @@ export function createProducedWorkItemInDatabase(
deliveryRole: input.delivery_role,
serviceId: input.service_id,
});
database
.prepare(
`INSERT INTO work_items (
group_folder,
chat_jid,
agent_type,
service_id,
delivery_role,
status,
start_seq,
end_seq,
result_payload,
attachment_payload,
delivery_attempts,
created_at,
updated_at
) VALUES (?, ?, ?, ?, ?, 'produced', ?, ?, ?, ?, 0, ?, ?)`,
)
.run(
input.group_folder,
input.chat_jid,
return database.transaction(() => {
supersedeOpenWorkItemsForCanonicalKey(
database,
input,
agentType,
serviceId,
input.delivery_role ?? null,
input.start_seq,
input.end_seq,
input.result_payload,
serializeAttachmentPayload(input.attachments),
now,
now,
);
const lastId = (
database.prepare('SELECT last_insert_rowid() as id').get() as { id: number }
).id;
return hydrateWorkItemRow(
database
.prepare('SELECT * FROM work_items WHERE id = ?')
.get(lastId) as StoredWorkItemRow,
);
.prepare(
`INSERT INTO work_items (
group_folder,
chat_jid,
agent_type,
service_id,
delivery_role,
status,
start_seq,
end_seq,
result_payload,
attachment_payload,
delivery_attempts,
created_at,
updated_at
) VALUES (?, ?, ?, ?, ?, 'produced', ?, ?, ?, ?, 0, ?, ?)`,
)
.run(
input.group_folder,
input.chat_jid,
agentType,
serviceId,
input.delivery_role ?? null,
input.start_seq,
input.end_seq,
input.result_payload,
serializeAttachmentPayload(input.attachments),
now,
now,
);
const lastId = (
database.prepare('SELECT last_insert_rowid() as id').get() as {
id: number;
}
).id;
return hydrateWorkItemRow(
database
.prepare('SELECT * FROM work_items WHERE id = ?')
.get(lastId) as StoredWorkItemRow,
);
})();
}
export function markWorkItemDeliveredInDatabase(