Skip to content
Merged
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
95 changes: 95 additions & 0 deletions src/__tests__/cli/commands/sync.skip-running.test.ts
Original file line number Diff line number Diff line change
@@ -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<typeof import('../../../config.js')>('../../../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);
});
});
92 changes: 92 additions & 0 deletions src/__tests__/cli/lock.test.ts
Original file line number Diff line number Diff line change
@@ -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<void> {
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<void> {
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);
});
});
61 changes: 61 additions & 0 deletions src/__tests__/utils/atomic.test.ts
Original file line number Diff line number Diff line change
@@ -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([]);
});
});
5 changes: 3 additions & 2 deletions src/cli/commands/init.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -113,13 +114,13 @@ export async function init(): Promise<InitResult> {
.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) {
Expand Down
16 changes: 15 additions & 1 deletion src/cli/commands/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -28,9 +29,14 @@ export interface OSRegistration {
*/
export async function sync(
register: (task: ScheduledTask, shimPath: string) => Promise<void>,
options?: { taskId?: string; configPath?: string },
options?: {
taskId?: string;
configPath?: string;
isTaskRunning?: (taskId: string) => Promise<boolean>;
},
): Promise<SyncResult> {
const configPath = options?.configPath ?? getGlobalSchedulesPath();
const isTaskRunning = options?.isTaskRunning ?? defaultIsTaskRunning;
let config = await loadConfig(configPath);

// Ensure executor is installed
Expand Down Expand Up @@ -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);
Expand Down
21 changes: 21 additions & 0 deletions src/cli/lock.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,27 @@ async function verifyProcessIdentity(pid: number, taskId: string): Promise<boole
}
}

/**
* Report whether a live, identity-verified executor currently holds this
* task's lock. Used by sync to avoid unloading/reloading a task's OS
* registration while its job is mid-run — which would tear down the running
* claude process. A dead or unverifiable lock reports false so a stale lock
* never blocks re-registration indefinitely.
*/
export async function isTaskRunning(taskId: string): Promise<boolean> {
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<boolean> {
const lockPath = getLockPath(taskId);
const pid = await readLockPid(taskId);
Expand Down
5 changes: 3 additions & 2 deletions src/cli/platform.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -78,7 +79,7 @@ async function registerDarwin(task: ScheduledTask, shimPath: string): Promise<vo
await defaultExec('launchctl', ['unload', plistPath]);
} catch { /* not loaded */ }

await fs.writeFile(plistPath, plistContent, 'utf-8');
await writeFileAtomic(plistPath, plistContent);
await defaultExec('launchctl', ['load', plistPath]);

// Ensure daily auto-sync job exists when registering timezone-aware tasks
Expand Down Expand Up @@ -145,7 +146,7 @@ async function ensureTzSyncJob(): Promise<void> {
await defaultExec('launchctl', ['unload', plistPath]);
} catch { /* not loaded */ }

await fs.writeFile(plistPath, plist, 'utf-8');
await writeFileAtomic(plistPath, plist);
await defaultExec('launchctl', ['load', plistPath]);
}

Expand Down
Loading
Loading