diff --git a/src/__tests__/cli/commands/sync.skip-running.test.ts b/src/__tests__/cli/commands/sync.skip-running.test.ts new file mode 100644 index 0000000..64d358f --- /dev/null +++ b/src/__tests__/cli/commands/sync.skip-running.test.ts @@ -0,0 +1,95 @@ +/** + * Tests for sync: skip re-registering tasks whose job is currently running. + */ + +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import { createTask, createEmptyConfig, addTask, type SchedulesConfig } from '../../../index.js'; + +const configStore = vi.hoisted((): { current: SchedulesConfig | null } => ({ current: null })); + +vi.mock('../../../cli/platform.js', () => ({ + registerTask: vi.fn(), + unregisterTask: vi.fn(), +})); + +vi.mock('../../../cli/commands/init.js', () => ({ + ensureExecutorInstalled: vi.fn(), + getShimPath: vi.fn(), + resolveClaudeBin: vi.fn(), +})); + +vi.mock('../../../config.js', async () => { + const actual = await vi.importActual('../../../config.js'); + return { + ...actual, + getGlobalSchedulesPath: vi.fn(), + getLogsDir: vi.fn(), + loadConfig: vi.fn(), + saveConfig: vi.fn(), + }; +}); + +import { sync } from '../../../cli/commands/sync.js'; +import * as initMod from '../../../cli/commands/init.js'; +import * as configMod from '../../../config.js'; + +function makeCronTask(name: string) { + return createTask({ + name, + trigger: { type: 'cron', expression: '0 7 * * *', timezone: 'UTC' }, + execution: { command: 'echo hi', workingDirectory: '/tmp' }, + }); +} + +describe('sync: skip tasks with a running job', () => { + beforeEach(() => { + configStore.current = createEmptyConfig(); + vi.mocked(initMod.ensureExecutorInstalled).mockResolvedValue({ + success: true, + executorPath: '/fake/executor', + shimPath: '/fake/shim', + cliShimPath: '/fake/cli-shim', + }); + vi.mocked(initMod.getShimPath).mockReturnValue('/fake/shim'); + vi.mocked(initMod.resolveClaudeBin).mockResolvedValue(undefined); + vi.mocked(configMod.getGlobalSchedulesPath).mockReturnValue('/fake/config.json'); + vi.mocked(configMod.getLogsDir).mockReturnValue('/tmp'); + vi.mocked(configMod.loadConfig).mockImplementation(async () => configStore.current as SchedulesConfig); + vi.mocked(configMod.saveConfig).mockImplementation(async (_path, cfg) => { + configStore.current = cfg as SchedulesConfig; + }); + }); + + it('does not re-register a task whose job is running, and syncs the rest', async () => { + const running = makeCronTask('running-task'); + const idle = makeCronTask('idle-task'); + configStore.current = addTask(addTask(createEmptyConfig(), running), idle); + + const register = vi.fn().mockResolvedValue(undefined); + const isTaskRunning = vi.fn(async (taskId: string) => taskId === running.id); + + const result = await sync(register, { isTaskRunning }); + + expect(result.success).toBe(true); + expect(result.synced).toEqual([idle.id]); + expect(result.skipped).toEqual([running.id]); + expect(register).toHaveBeenCalledTimes(1); + expect(register).toHaveBeenCalledWith(expect.objectContaining({ id: idle.id }), '/fake/shim'); + expect(register).not.toHaveBeenCalledWith(expect.objectContaining({ id: running.id }), expect.anything()); + }); + + it('registers all tasks when none are running', async () => { + const a = makeCronTask('a'); + const b = makeCronTask('b'); + configStore.current = addTask(addTask(createEmptyConfig(), a), b); + + const register = vi.fn().mockResolvedValue(undefined); + const isTaskRunning = vi.fn().mockResolvedValue(false); + + const result = await sync(register, { isTaskRunning }); + + expect(result.synced).toEqual([a.id, b.id]); + expect(result.skipped).toEqual([]); + expect(register).toHaveBeenCalledTimes(2); + }); +}); diff --git a/src/__tests__/cli/lock.test.ts b/src/__tests__/cli/lock.test.ts new file mode 100644 index 0000000..618620b --- /dev/null +++ b/src/__tests__/cli/lock.test.ts @@ -0,0 +1,92 @@ +/** + * Tests for isTaskRunning: does a live, identity-verified executor hold the lock? + */ + +import { describe, it, expect, afterEach } from 'vitest'; +import { mkdir, writeFile, rm } from 'node:fs/promises'; +import { spawn, type ChildProcess } from 'node:child_process'; +import path from 'node:path'; +import { getLockPath, isTaskRunning } from '../../cli/lock.js'; + +let counter = 0; +const created: string[] = []; +const children: ChildProcess[] = []; + +function uniqueTaskId(): string { + const id = `lock-test-${process.pid}-${counter++}`; + created.push(getLockPath(id)); + return id; +} + +async function writeLock(taskId: string, pid: number, startTime?: number): Promise { + const dir = getLockPath(taskId); + await mkdir(dir, { recursive: true }); + if (startTime !== undefined) { + await writeFile(path.join(dir, 'startTime'), String(startTime), 'utf-8'); + } + await writeFile(path.join(dir, 'pid'), String(pid), 'utf-8'); +} + +function spawnLiveChild(): ChildProcess { + const child = spawn(process.execPath, ['-e', 'setInterval(() => {}, 1000)'], { stdio: 'ignore' }); + children.push(child); + return child; +} + +async function waitUntilDead(pid: number): Promise { + const deadline = Date.now() + 5000; + while (Date.now() < deadline) { + try { + process.kill(pid, 0); + } catch { + return; + } + await new Promise(r => setTimeout(r, 50)); + } +} + +afterEach(async () => { + for (const child of children.splice(0)) { + if (child.pid && !child.killed) child.kill('SIGKILL'); + } + for (const dir of created.splice(0)) { + await rm(dir, { recursive: true, force: true }).catch(() => {}); + } +}); + +describe('isTaskRunning', () => { + it('returns false when no lock exists', async () => { + expect(await isTaskRunning(uniqueTaskId())).toBe(false); + }); + + it('returns false when the lock has no startTime (unverifiable)', async () => { + const taskId = uniqueTaskId(); + await writeLock(taskId, process.pid); // alive pid, but no startTime + expect(await isTaskRunning(taskId)).toBe(false); + }); + + it('returns false when the locked process is dead', async () => { + const taskId = uniqueTaskId(); + const child = spawnLiveChild(); + const pid = child.pid!; + const startTime = Date.now(); + child.kill('SIGKILL'); + await waitUntilDead(pid); + await writeLock(taskId, pid, startTime); + expect(await isTaskRunning(taskId)).toBe(false); + }); + + it('returns true for a live, identity-verified process', async () => { + const taskId = uniqueTaskId(); + const child = spawnLiveChild(); + const startTime = Date.now(); + await new Promise(r => setTimeout(r, 100)); // let ps see the process + await writeLock(taskId, child.pid!, startTime); + + expect(await isTaskRunning(taskId)).toBe(true); + + child.kill('SIGKILL'); + await waitUntilDead(child.pid!); + expect(await isTaskRunning(taskId)).toBe(false); + }); +}); diff --git a/src/__tests__/utils/atomic.test.ts b/src/__tests__/utils/atomic.test.ts new file mode 100644 index 0000000..77e7356 --- /dev/null +++ b/src/__tests__/utils/atomic.test.ts @@ -0,0 +1,61 @@ +import { describe, it, expect, beforeEach, afterEach } from 'vitest'; +import { mkdir, rm, readFile, writeFile, stat, readdir } from 'node:fs/promises'; +import path from 'node:path'; +import os from 'node:os'; +import { writeFileAtomic } from '../../utils/atomic.js'; + +describe('writeFileAtomic', () => { + const tmpDir = path.join(os.tmpdir(), `atomic-test-${process.pid}`); + + beforeEach(async () => { + await mkdir(tmpDir, { recursive: true }); + }); + + afterEach(async () => { + await rm(tmpDir, { recursive: true, force: true }); + }); + + it('writes content to a new file', async () => { + const target = path.join(tmpDir, 'new.txt'); + await writeFileAtomic(target, 'hello'); + expect(await readFile(target, 'utf-8')).toBe('hello'); + }); + + it('overwrites an existing file', async () => { + const target = path.join(tmpDir, 'existing.txt'); + await writeFile(target, 'old', 'utf-8'); + await writeFileAtomic(target, 'new'); + expect(await readFile(target, 'utf-8')).toBe('new'); + }); + + it('applies the requested mode', async () => { + const target = path.join(tmpDir, 'exec.sh'); + await writeFileAtomic(target, '#!/bin/bash\n', { mode: 0o755 }); + const mode = (await stat(target)).mode & 0o777; + expect(mode).toBe(0o755); + }); + + it('leaves no temp file behind on success', async () => { + const target = path.join(tmpDir, 'clean.txt'); + await writeFileAtomic(target, 'data'); + const leftovers = (await readdir(tmpDir)).filter(f => f.includes('.tmp-')); + expect(leftovers).toEqual([]); + }); + + it('handles concurrent writes to the same target without colliding', async () => { + const target = path.join(tmpDir, 'concurrent.txt'); + await Promise.all( + Array.from({ length: 10 }, () => writeFileAtomic(target, 'same')), + ); + expect(await readFile(target, 'utf-8')).toBe('same'); + const leftovers = (await readdir(tmpDir)).filter(f => f.includes('.tmp-')); + expect(leftovers).toEqual([]); + }); + + it('cleans up the temp file and rejects when the target dir is missing', async () => { + const target = path.join(tmpDir, 'missing-subdir', 'file.txt'); + await expect(writeFileAtomic(target, 'data')).rejects.toThrow(); + const leftovers = (await readdir(tmpDir)).filter(f => f.includes('.tmp-')); + expect(leftovers).toEqual([]); + }); +}); diff --git a/src/cli/commands/init.ts b/src/cli/commands/init.ts index ac4a063..aec2a2d 100644 --- a/src/cli/commands/init.ts +++ b/src/cli/commands/init.ts @@ -15,6 +15,7 @@ import fs from 'node:fs/promises'; import { execFile } from 'node:child_process'; import path from 'node:path'; import os from 'node:os'; +import { writeFileAtomic } from '../../utils/atomic.js'; export interface InitResult { success: boolean; @@ -113,13 +114,13 @@ export async function init(): Promise { .replace('{{EXECUTOR_PATH}}', executorPath) .replace('{{NODE_PATH}}', nodePath) .replace('{{USER_PATH}}', userPath); - await fs.writeFile(shimDest, shimContent, { mode: 0o755 }); + await writeFileAtomic(shimDest, shimContent, { mode: 0o755 }); // Write CLI shim (no PATH restore needed — runs interactively) const cliShimContent = CLI_SHIM_TEMPLATE .replace('{{CLI_ENTRY_PATH}}', cliEntryPath) .replace('{{NODE_PATH}}', nodePath); - await fs.writeFile(cliShimDest, cliShimContent, { mode: 0o755 }); + await writeFileAtomic(cliShimDest, cliShimContent, { mode: 0o755 }); return { success: true, executorPath, shimPath: shimDest, cliShimPath: cliShimDest }; } catch (err) { diff --git a/src/cli/commands/sync.ts b/src/cli/commands/sync.ts index 5e79d14..2db4759 100644 --- a/src/cli/commands/sync.ts +++ b/src/cli/commands/sync.ts @@ -8,6 +8,7 @@ import path from 'node:path'; import { loadConfig, saveConfig, updateTask, getGlobalSchedulesPath, getLogsDir } from '../../config.js'; import { getShimPath, ensureExecutorInstalled, resolveClaudeBin } from './init.js'; import { unregisterTask } from '../platform.js'; +import { isTaskRunning as defaultIsTaskRunning } from '../lock.js'; import type { ScheduledTask } from '../../types.js'; export interface SyncResult { @@ -28,9 +29,14 @@ export interface OSRegistration { */ export async function sync( register: (task: ScheduledTask, shimPath: string) => Promise, - options?: { taskId?: string; configPath?: string }, + options?: { + taskId?: string; + configPath?: string; + isTaskRunning?: (taskId: string) => Promise; + }, ): Promise { const configPath = options?.configPath ?? getGlobalSchedulesPath(); + const isTaskRunning = options?.isTaskRunning ?? defaultIsTaskRunning; let config = await loadConfig(configPath); // Ensure executor is installed @@ -71,6 +77,14 @@ export async function sync( continue; } + // Skip tasks whose job is mid-run: re-registering unloads/reloads the + // launchd service and would tear down the running claude process. The + // next sync re-registers it once the run finishes. + if (await isTaskRunning(task.id)) { + skipped.push(task.id); + continue; + } + try { await register(task, shimPath); synced.push(task.id); diff --git a/src/cli/lock.ts b/src/cli/lock.ts index d3b5061..b4d652a 100644 --- a/src/cli/lock.ts +++ b/src/cli/lock.ts @@ -56,6 +56,27 @@ async function verifyProcessIdentity(pid: number, taskId: string): Promise { + const pid = await readLockPid(taskId); + if (pid === null) return false; + + if (!(await verifyProcessIdentity(pid, taskId))) return false; + + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + export async function killRunningTask(taskId: string): Promise { const lockPath = getLockPath(taskId); const pid = await readLockPid(taskId); diff --git a/src/cli/platform.ts b/src/cli/platform.ts index df1c52a..f883150 100644 --- a/src/cli/platform.ts +++ b/src/cli/platform.ts @@ -7,6 +7,7 @@ import fs from 'node:fs/promises'; import path from 'node:path'; import os from 'node:os'; import { exec as defaultExec } from '../utils/exec.js'; +import { writeFileAtomic } from '../utils/atomic.js'; import { getLogsDir } from '../config.js'; import { generatePlist, @@ -78,7 +79,7 @@ async function registerDarwin(task: ScheduledTask, shimPath: string): Promise { await defaultExec('launchctl', ['unload', plistPath]); } catch { /* not loaded */ } - await fs.writeFile(plistPath, plist, 'utf-8'); + await writeFileAtomic(plistPath, plist); await defaultExec('launchctl', ['load', plistPath]); } diff --git a/src/utils/atomic.ts b/src/utils/atomic.ts new file mode 100644 index 0000000..f56c5a1 --- /dev/null +++ b/src/utils/atomic.ts @@ -0,0 +1,38 @@ +/** + * Atomic file writes: write to a temp file in the same directory, then rename + * over the target. A concurrent reader (e.g. launchd exec'ing a shim while a + * sync rewrites it) always sees either the whole old file or the whole new + * one, never a truncated or half-written file. + */ + +import fs from 'node:fs/promises'; +import path from 'node:path'; + +export interface WriteFileAtomicOptions { + mode?: number; +} + +// Distinguishes concurrent writes to the same target from within one process, +// so their temp files never collide. +let writeCounter = 0; + +export async function writeFileAtomic( + filePath: string, + content: string, + options?: WriteFileAtomicOptions, +): Promise { + const dir = path.dirname(filePath); + const tmpPath = path.join(dir, `.${path.basename(filePath)}.tmp-${process.pid}-${writeCounter++}`); + + const writeOptions = options?.mode !== undefined ? { mode: options.mode } : undefined; + + try { + await fs.writeFile(tmpPath, content, writeOptions); + // rename is atomic on the same filesystem; the temp file shares the target's + // directory so this never crosses a filesystem boundary. + await fs.rename(tmpPath, filePath); + } catch (err) { + await fs.rm(tmpPath, { force: true }).catch(() => {}); + throw err; + } +}