File size: 4,367 Bytes
5710dd0
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
// apps/vis/server/src/lib/cron-store.ts
//
// Read-only reader for cron tasks. v2 persists cron state as durable wire
// records (`cron.add` / `cron.delete` / `cron.cursor`) inside each agent's
// `wire.jsonl` (`packages/agent-core-v2/src/features/cron/`); v1 wrote
// `<agentDir>/cron/<id>.json` files instead. vis reads both: legacy files
// first, then the wire fold on top, so a session written by either engine
// (or carried across both) lists its cron tasks. The visualizer never
// writes anything.

import { readdir, readFile } from 'node:fs/promises';
import { join } from 'node:path';

import type { CronTask } from './agent-record-types';
import { readAgentWire } from './wire-reader';

/** Cron id format: 8 lowercase hex chars (legacy v1 ids) or a 26-char ULID
 *  (v2 ids) — mirror of the engine's `CRON_ID_REGEX`
 *  (`features/cron/cronService.ts`). */
const VALID_CRON_ID = /^(?:[0-9a-f]{8}|[0-9A-HJKMNP-TV-Z]{26})$/i;

export function isSafeCronId(id: string): boolean {
  return VALID_CRON_ID.test(id);
}

function cronDirOf(agentDir: string): string {
  return join(agentDir, 'cron');
}

/**
 * Enumerate all cron tasks for one agent homedir, sorted by creation time
 * (oldest first, matching how a user scheduled them).
 *
 * Legacy `<agentDir>/cron/*.json` files whose names don't match
 * `VALID_CRON_ID`, fail to parse, or miss required fields are skipped;
 * a missing/unreadable `wire.jsonl` contributes no records.
 */
export async function listCronTasks(agentDir: string): Promise<CronTask[]> {
  const byId = new Map<string, CronTask>();
  for (const task of await listCronTaskFiles(agentDir)) {
    byId.set(task.id, task);
  }
  await foldCronWireInto(agentDir, byId);
  return [...byId.values()].sort((a, b) => a.createdAt - b.createdAt);
}

/** Legacy v1 layout: one JSON file per task under `<agentDir>/cron/`. */
async function listCronTaskFiles(agentDir: string): Promise<CronTask[]> {
  const dir = cronDirOf(agentDir);
  let entries: import('node:fs').Dirent[];
  try {
    entries = await readdir(dir, { withFileTypes: true });
  } catch {
    return [];
  }
  const out: CronTask[] = [];
  for (const entry of entries) {
    if (!entry.isFile() || !entry.name.endsWith('.json')) continue;
    const id = entry.name.slice(0, -'.json'.length);
    if (!VALID_CRON_ID.test(id)) continue;
    let parsed: unknown;
    try {
      parsed = JSON.parse(await readFile(join(dir, entry.name), 'utf8'));
    } catch {
      continue;
    }
    if (isCronTask(parsed)) out.push(parsed);
  }
  return out;
}

/** v2 layout: fold the agent's wire records — `cron.add` upserts,
 *  `cron.delete` removes, `cron.cursor` advances `lastFiredAt`. The fold
 *  applies on top of the legacy-file state, so the wire (the engine's
 *  authoritative journal) wins for tasks present in both. */
async function foldCronWireInto(agentDir: string, byId: Map<string, CronTask>): Promise<void> {
  let records;
  try {
    ({ records } = await readAgentWire(join(agentDir, 'wire.jsonl')));
  } catch {
    return;
  }
  for (const entry of records) {
    const rec = entry.data;
    switch (rec.type) {
      case 'cron.add':
        if (isCronTask(rec.task)) byId.set(rec.task.id, rec.task);
        break;
      case 'cron.delete': {
        // Tolerate hand-edited / partially corrupted wires: the reader only
        // validates the record's `type`, so guard the payload shape before
        // iterating, the same way `cron.add` goes through `isCronTask`.
        const ids: unknown = rec.ids;
        if (Array.isArray(ids)) {
          for (const id of ids) if (typeof id === 'string') byId.delete(id);
        }
        break;
      }
      case 'cron.cursor': {
        const id: unknown = rec.id;
        const lastFiredAt: unknown = rec.lastFiredAt;
        if (typeof id !== 'string' || typeof lastFiredAt !== 'number') break;
        const task = byId.get(id);
        if (task !== undefined) byId.set(id, { ...task, lastFiredAt });
        break;
      }
      default:
        break;
    }
  }
}

function isCronTask(value: unknown): value is CronTask {
  if (typeof value !== 'object' || value === null) return false;
  const o = value as Record<string, unknown>;
  return (
    typeof o['id'] === 'string' &&
    typeof o['cron'] === 'string' &&
    typeof o['prompt'] === 'string' &&
    typeof o['createdAt'] === 'number'
  );
}