summaryrefslogtreecommitdiffhomepage
path: root/packages/heartbeat/src/run-store.ts
blob: 9650181582a360eb7003c627ecc9a7754d5de977 (plain)
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
import type { StorageNamespace } from "@dispatch/kernel";
import type { HeartbeatRun, HeartbeatRunStatus } from "@dispatch/transport-contract";

/** Storage key for a single heartbeat run record. */
function runKey(workspaceId: string, runId: string): string {
  return `run:${workspaceId}:${runId}`;
}

/** Prefix matching every run record for a workspace (for enumeration). */
function runPrefix(workspaceId: string): string {
  return `run:${workspaceId}:`;
}

/** Extract the runId from a full `run:<workspaceId>:<runId>` key. */
function parseRunId(key: string, workspaceId: string): string {
  const prefix = runPrefix(workspaceId);
  return key.startsWith(prefix) ? key.slice(prefix.length) : key;
}

export interface HeartbeatRunStore {
  /** Create a new run record (status `"running"`) and persist it. */
  readonly create: (workspaceId: string, run: HeartbeatRun) => Promise<HeartbeatRun>;
  /** Update the status of an existing run. No-op if the run is unknown. */
  readonly setStatus: (
    workspaceId: string,
    runId: string,
    status: HeartbeatRunStatus,
  ) => Promise<HeartbeatRun | null>;
  /** A single run by id, or `null` when unknown. */
  readonly get: (workspaceId: string, runId: string) => Promise<HeartbeatRun | null>;
  /** All runs for a workspace, most-recent first (by `triggeredAt`). */
  readonly list: (workspaceId: string) => Promise<readonly HeartbeatRun[]>;
}

export function createHeartbeatRunStore(storage: StorageNamespace): HeartbeatRunStore {
  async function readRun(workspaceId: string, runId: string): Promise<HeartbeatRun | null> {
    const raw = await storage.get(runKey(workspaceId, runId));
    if (raw === null) return null;
    try {
      const parsed = JSON.parse(raw) as Partial<HeartbeatRun>;
      if (
        typeof parsed.id !== "string" ||
        typeof parsed.conversationId !== "string" ||
        typeof parsed.triggeredAt !== "string" ||
        typeof parsed.status !== "string"
      ) {
        return null;
      }
      return {
        id: parsed.id,
        conversationId: parsed.conversationId,
        triggeredAt: parsed.triggeredAt,
        status: parsed.status as HeartbeatRunStatus,
      };
    } catch {
      return null;
    }
  }

  return {
    async create(workspaceId: string, run: HeartbeatRun): Promise<HeartbeatRun> {
      await storage.set(runKey(workspaceId, run.id), JSON.stringify(run));
      return run;
    },

    async setStatus(
      workspaceId: string,
      runId: string,
      status: HeartbeatRunStatus,
    ): Promise<HeartbeatRun | null> {
      const run = await readRun(workspaceId, runId);
      if (run === null) return null;
      const updated: HeartbeatRun = { ...run, status };
      await storage.set(runKey(workspaceId, runId), JSON.stringify(updated));
      return updated;
    },

    async get(workspaceId: string, runId: string): Promise<HeartbeatRun | null> {
      return readRun(workspaceId, runId);
    },

    async list(workspaceId: string): Promise<readonly HeartbeatRun[]> {
      const keys = await storage.keys(runPrefix(workspaceId));
      const runs: HeartbeatRun[] = [];
      for (const key of keys) {
        const runId = parseRunId(key, workspaceId);
        const run = await readRun(workspaceId, runId);
        if (run !== null) runs.push(run);
      }
      // Most-recent first by triggeredAt (ISO-8601 sorts lexicographically).
      runs.sort((a, b) => b.triggeredAt.localeCompare(a.triggeredAt));
      return runs;
    },
  };
}