Skip to content
19 changes: 19 additions & 0 deletions backend/src/db/migrations/022-schedule-runs-session-index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
import type { Migration } from '../migration-runner'

const migration: Migration = {
version: 22,
name: 'schedule-runs-session-index',

up(db) {
db.run(`
CREATE INDEX IF NOT EXISTS idx_schedule_runs_session
ON schedule_runs(session_id, started_at DESC)
`)
},

down(db) {
db.run('DROP INDEX IF EXISTS idx_schedule_runs_session')
},
}

export default migration
2 changes: 2 additions & 0 deletions backend/src/db/migrations/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import migration018 from './018-session-pins'
import migration019 from './019-drop-opencode-configs'
import migration020 from './020-drop-opencode-model-state'
import migration021 from './021-drop-schedule-run-workspace-id'
import migration022 from './022-schedule-runs-session-index'

export const allMigrations: Migration[] = [
migration001,
Expand All @@ -43,4 +44,5 @@ export const allMigrations: Migration[] = [
migration019,
migration020,
migration021,
migration022,
]
13 changes: 12 additions & 1 deletion backend/src/db/schedules.ts
Original file line number Diff line number Diff line change
Expand Up @@ -460,6 +460,12 @@ export function getScheduleRunById(db: Database, repoId: number, jobId: number,
return row ? rowToScheduleRun(row) : null
}

export function getScheduleRunBySessionId(db: Database, sessionId: string): ScheduleRun | null {
const stmt = db.prepare('SELECT * FROM schedule_runs WHERE session_id = ? ORDER BY started_at DESC LIMIT 1')
const row = stmt.get(sessionId) as ScheduleRunRow | undefined
return row ? rowToScheduleRun(row) : null
}

export function getRunningScheduleRunByJob(db: Database, repoId: number, jobId: number): ScheduleRun | null {
const stmt = db.prepare(`
SELECT * FROM schedule_runs
Expand Down Expand Up @@ -592,10 +598,11 @@ export interface ListAllRunsOptions {
repoId?: number
jobId?: number
triggerSource?: string
runId?: number
}

export function listAllScheduleRuns(db: Database, options: ListAllRunsOptions = {}): ScheduleRunWithContext[] {
const { limit = 50, offset = 0, status, repoId, jobId, triggerSource } = options
const { limit = 50, offset = 0, status, repoId, jobId, triggerSource, runId } = options
const conditions: string[] = []
const params: (string | number)[] = []

Expand All @@ -615,6 +622,10 @@ export function listAllScheduleRuns(db: Database, options: ListAllRunsOptions =
conditions.push('sr.trigger_source = ?')
params.push(triggerSource)
}
if (runId !== undefined) {
conditions.push('sr.id = ?')
params.push(runId)
}

const whereClause = conditions.length > 0 ? `WHERE ${conditions.join(' AND ')}` : ''

Expand Down
2 changes: 1 addition & 1 deletion backend/src/routes/internal/notifications.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ export function createInternalNotificationRoutes(notificationService: Notificati
timestamp: Date.now(),
data: {
eventType: 'assistant.message',
url: parsed.data.url ?? '/',
url: notificationService.getScheduleRunReportUrl(parsed.data.sessionId) ?? parsed.data.url ?? '/',
priority: parsed.data.priority,
},
}
Expand Down
20 changes: 19 additions & 1 deletion backend/src/routes/schedules.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,23 @@ function parseRunListLimit(value: string | undefined): number {
return Math.min(parsed, 100)
}

function parseRunIdFilter(value: string | undefined): number | undefined {
if (value === undefined) {
return undefined
}

if (!/^\d+$/.test(value)) {
throw new ScheduleServiceError('Invalid run id', 400)
}

const parsed = Number(value)
if (!Number.isSafeInteger(parsed) || parsed < 1) {
throw new ScheduleServiceError('Run id must be a positive integer', 400)
}

return parsed
}

export function createScheduleRoutes(scheduleService: ScheduleService) {
const app = new Hono()

Expand Down Expand Up @@ -48,7 +65,8 @@ export function createScheduleRoutes(scheduleService: ScheduleService) {
return Number.isNaN(parsed) ? undefined : parsed
})() : undefined
const triggerSource = c.req.query('triggerSource') || undefined
const runs = scheduleService.listAllRuns({ limit, offset, status, repoId, jobId, triggerSource })
const runId = parseRunIdFilter(c.req.query('runId'))
const runs = scheduleService.listAllRuns({ limit, offset, status, repoId, jobId, triggerSource, runId })
return c.json({ runs })
} catch (error) {
return handleServiceError(c, error, 'Failed to list all schedule runs', ScheduleServiceError)
Expand Down
1 change: 1 addition & 0 deletions backend/src/services/assistant-mode.ts
Original file line number Diff line number Diff line change
Expand Up @@ -569,6 +569,7 @@ Sending is rate limited to **10 notifications per minute**. Beyond that the tool
- Notifications are only sent if the user has registered devices (browser push subscriptions)
- If VAPID is not configured on the server, the tool fails with a \`503\` status
- Use \`priority: 'high'\` for urgent notifications that should interrupt the user
- In a scheduled run, the notification always opens that run's report, so \`url\` is not needed
- Do not call the internal HTTP API with \`curl\` for notifications; the tool is the supported path
`
}
Expand Down
20 changes: 19 additions & 1 deletion backend/src/services/notification.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import {
getRepoName,
listRepos,
} from "../db/queries";
import { getScheduleRunBySessionId } from "../db/schedules";
import type { Repo } from "../types/repo";
import { getReposPath } from "@opencode-manager/shared/config/env";
import { ASSISTANT_REPO_ID } from "@opencode-manager/shared/utils";
Expand Down Expand Up @@ -71,13 +72,22 @@ const EVENT_CONFIG: Record<

const MAX_BODY_LENGTH = 140;

const RUN_OUTCOME_EVENTS = new Set<string>([
NotificationEventType.SESSION_IDLE,
NotificationEventType.SESSION_FAILED,
]);

function resolveEventSessionId(event: SSEEvent): string | undefined {
if (event.type === NotificationEventType.FORM_CREATED) {
return event.data.form.sessionID;
}
return sessionIDFromEvent(event);
}

export function buildScheduleRunReportUrl(runId: number): string {
return `/schedules?scheduleTab=runs&runId=${runId}`;
}

export function buildNotificationUrl(
repo: Pick<Repo, "id"> | null,
sessionId: string | undefined
Expand Down Expand Up @@ -272,6 +282,13 @@ export class NotificationService {
return rows.map((r) => r.user_id);
}

/** Returns the run report URL when the session was started by a scheduled run. */
getScheduleRunReportUrl(sessionId: string | undefined): string | null {
if (!sessionId) return null;
const run = getScheduleRunBySessionId(this.db, sessionId);
return run ? buildScheduleRunReportUrl(run.id) : null;
}

private async resolveRepoForDirectory(directory: string): Promise<Repo | null> {
const repo =
getRepoBySourcePath(this.db, path.resolve(directory)) ??
Expand Down Expand Up @@ -312,7 +329,8 @@ export class NotificationService {
const repo = directory ? await this.resolveRepoForDirectory(directory) : null;
const repoId = repo?.id;
const repoName = repo ? getRepoName(repo) : undefined;
const url = buildNotificationUrl(repo, sessionId);
const reportUrl = RUN_OUTCOME_EVENTS.has(event.type) ? this.getScheduleRunReportUrl(sessionId) : null;
const url = reportUrl ?? buildNotificationUrl(repo, sessionId);

const payload = buildEventNotificationPayload(event, {
repoName,
Expand Down
15 changes: 8 additions & 7 deletions backend/src/services/opencode-manager-tool-plugin.ts
Original file line number Diff line number Diff line change
Expand Up @@ -198,18 +198,19 @@ async function postInternalApi(routePath, body, signal) {

var ACTIONS = {
send_notification: {
run: async function (params, signal) {
var result = await postInternalApi('/notifications/send', params, signal)
run: async function (params, context) {
var body = Object.assign({}, params, { sessionId: context.sessionID })
var result = await postInternalApi('/notifications/send', body, context.signal)
if (result.noSubscriptions === true) {
return 'No devices are registered for push notifications, so nothing was delivered.'
}
return 'Notification sent: ' + (result.delivered || 0) + ' delivered, ' + (result.failed || 0) + ' failed.'
},
},
request: {
run: async function (params, signal) {
run: async function (params, context) {
assertAllowedRoute(params.method, params.path)
var text = await requestInternalApi(params.method, params.path, params.body, signal)
var text = await requestInternalApi(params.method, params.path, params.body, context.signal)
return text || 'The request succeeded with an empty response body.'
},
},
Expand All @@ -229,14 +230,14 @@ function assertParams(actionName, params) {
}
}

async function runAction(input, signal) {
async function runAction(input, context) {
var actionName = input !== null && typeof input === 'object' ? input.action : undefined
if (!Object.prototype.hasOwnProperty.call(ACTIONS, actionName)) {
throw new Error('Unknown OpenCode Manager action: ' + String(actionName) + '. Supported actions: ' + ACTION_NAMES.join(', ') + '.')
}
var params = input.params
assertParams(actionName, params)
return await ACTIONS[actionName].run(params, signal)
return await ACTIONS[actionName].run(params, context)
}

export default {
Expand All @@ -249,7 +250,7 @@ export default {
input: INPUT_SCHEMA,
options: { codemode: false },
execute: async function (input, context) {
return { content: await runAction(input, context.signal) }
return { content: await runAction(input, context) }
},
})
})
Expand Down
31 changes: 31 additions & 0 deletions backend/test/db/schedule-migrations.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import migration007 from '../../src/db/migrations/007-schedules'
import migration008 from '../../src/db/migrations/008-schedule-cron-support'
import migration015 from '../../src/db/migrations/015-schedule-worktree-isolation'
import migration021 from '../../src/db/migrations/021-drop-schedule-run-workspace-id'
import migration022 from '../../src/db/migrations/022-schedule-runs-session-index'

describe('schedule migrations', () => {
it('creates schedule jobs with nullable interval minutes in v7', () => {
Expand Down Expand Up @@ -206,3 +207,33 @@ describe('migration 021 - drop schedule run workspace id', () => {
db.close()
})
})

describe('migration 022 - schedule run session index', () => {
it('creates the composite session index on schedule_runs', () => {
const db = new Database(':memory:')
db.run('CREATE TABLE schedule_runs (id INTEGER PRIMARY KEY, session_id TEXT, started_at INTEGER NOT NULL)')

migration022.up(db)

const indexes = (db.prepare('PRAGMA index_list(schedule_runs)').all() as { name: string }[]).map((index) => index.name)
expect(indexes).toContain('idx_schedule_runs_session')

const columns = (db.prepare('PRAGMA index_info(idx_schedule_runs_session)').all() as { name: string }[]).map((column) => column.name)
expect(columns).toEqual(['session_id', 'started_at'])

db.close()
})

it('drops the index on rollback', () => {
const db = new Database(':memory:')
db.run('CREATE TABLE schedule_runs (id INTEGER PRIMARY KEY, session_id TEXT, started_at INTEGER NOT NULL)')

migration022.up(db)
migration022.down(db)

const indexes = (db.prepare('PRAGMA index_list(schedule_runs)').all() as { name: string }[]).map((index) => index.name)
expect(indexes).not.toContain('idx_schedule_runs_session')

db.close()
})
})
19 changes: 19 additions & 0 deletions backend/test/db/schedules.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -580,6 +580,25 @@ describe('schedule database queries', () => {
expect(stmt.all).toHaveBeenCalledWith('manual', 50, 0)
})

it('applies runId filter', () => {
const stmt = { all: vi.fn().mockReturnValue([]) }
mockDb.prepare.mockReturnValue(stmt)

schedulesDb.listAllScheduleRuns(mockDb, { runId: 99 })

expect(mockDb.prepare).toHaveBeenCalledWith(expect.stringContaining('sr.id = ?'))
expect(stmt.all).toHaveBeenCalledWith(99, 50, 0)
})

it('applies runId together with other filters', () => {
const stmt = { all: vi.fn().mockReturnValue([]) }
mockDb.prepare.mockReturnValue(stmt)

schedulesDb.listAllScheduleRuns(mockDb, { runId: 99, repoId: 42, limit: 1 })

expect(stmt.all).toHaveBeenCalledWith(42, 99, 1, 0)
})

it('applies limit and offset', () => {
const stmt = { all: vi.fn().mockReturnValue([]) }
mockDb.prepare.mockReturnValue(stmt)
Expand Down
37 changes: 37 additions & 0 deletions backend/test/routes/internal-notifications.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import { createOpenCodeClient } from '../../src/services/opencode/client'
import { allMigrations } from '../../src/db/migrations'
import { getOrCreateInternalToken } from '../../src/services/internal-token'
import { migrate } from '../../src/db/migration-runner'
import { createScheduleRun, updateScheduleRunMetadata } from '../../src/db/schedules'
import type { ScheduleWorktreeManager } from '../../src/services/schedule-worktree'

describe('internal/notifications routes', () => {
Expand Down Expand Up @@ -142,6 +143,42 @@ describe('internal/notifications routes', () => {
expect(res.status).toBe(400)
})

describe('notification url', () => {
const send = (payload: Record<string, unknown>) =>
app.request('/api/internal/notifications/send', {
method: 'POST',
body: JSON.stringify({ title: 'Test', body: 'Body', ...payload }),
headers: { 'content-type': 'application/json', authorization: `Bearer ${token}` },
})

const sentUrl = (sendToUser: ReturnType<typeof vi.spyOn>) =>
(sendToUser.mock.calls[0]?.[1] as { data: { url: string } }).data.url

beforeEach(() => {
vi.spyOn(notificationService, 'isConfigured').mockReturnValue(true)
})

it('links a scheduled run session to its run report, even when the agent passes a url', async () => {
db.exec('PRAGMA foreign_keys = OFF')
const run = createScheduleRun(db, { jobId: 7, repoId: 0, triggerSource: 'schedule', status: 'running', startedAt: 1, createdAt: 1 })
updateScheduleRunMetadata(db, 0, 7, run.id, { sessionId: 'ses_scheduled' })
const sendToUser = vi.spyOn(notificationService, 'sendToUser').mockResolvedValue({ delivered: 1, expired: 0, failed: 0, total: 1 })

const res = await send({ sessionId: 'ses_scheduled', url: '/repos/my-repo' })

expect(res.status).toBe(200)
expect(sentUrl(sendToUser)).toBe(`/schedules?scheduleTab=runs&runId=${run.id}`)
})

it('keeps the agent url for sessions that are not scheduled runs', async () => {
const sendToUser = vi.spyOn(notificationService, 'sendToUser').mockResolvedValue({ delivered: 1, expired: 0, failed: 0, total: 1 })

await send({ sessionId: 'ses_manual', url: '/repos/3' })

expect(sentUrl(sendToUser)).toBe('/repos/3')
})
})

it('POST /api/internal/notifications/send returns 429 after 10 calls within rate window', async () => {
vi.spyOn(notificationService, 'isConfigured').mockReturnValue(true)
vi.spyOn(notificationService, 'sendToUser').mockResolvedValue({ delivered: 0, expired: 0, failed: 0, total: 0 })
Expand Down
29 changes: 29 additions & 0 deletions backend/test/routes/schedules.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -273,4 +273,33 @@ describe('Schedule Routes', () => {
triggerSource: 'manual',
})
})

it('passes a runId bound to the service for all runs', async () => {
scheduleService.listAllRuns.mockReturnValue([])

const response = await app.request('/repos/42/schedules/all/runs?runId=99')
const body = await response.json() as { runs: Array<unknown> }

expect(response.status).toBe(200)
expect(body.runs).toHaveLength(0)
expect(scheduleService.listAllRuns).toHaveBeenCalledWith(expect.objectContaining({ runId: 99 }))
})

it.each(['abc', '1abc', '1.5', '-1', ''])('rejects an invalid runId format %s for all runs', async (runId) => {
const response = await app.request(`/repos/42/schedules/all/runs?runId=${encodeURIComponent(runId)}`)
const body = await response.json() as { error: string }

expect(response.status).toBe(400)
expect(body.error).toBe('Invalid run id')
expect(scheduleService.listAllRuns).not.toHaveBeenCalled()
})

it.each(['0', '99999999999999999999'])('rejects a non-positive or unsafe runId %s for all runs', async (runId) => {
const response = await app.request(`/repos/42/schedules/all/runs?runId=${runId}`)
const body = await response.json() as { error: string }

expect(response.status).toBe(400)
expect(body.error).toBe('Run id must be a positive integer')
expect(scheduleService.listAllRuns).not.toHaveBeenCalled()
})
})
Loading