automation-ingest-store.mjs
235 lines 6.9 KB
Raw
sha256:700fafdd1afa490919f9515d660ca6e75456bcd5bb67513abcd8757a634c01f6 docs: record AIP-b SD-21 land (KN #308) Human 9 days ago
1 /**
2 * AIP — ingest rules + idempotency store (file + dedicated Netlify blobs).
3 * D26 / D27. Never reuse gateway-agent-credentials or gateway-billing.
4 */
5
6 import fs from 'node:fs/promises';
7 import path from 'node:path';
8 import { fileURLToPath } from 'node:url';
9 import { listPackTemplates } from '../../lib/automation-ingest-policy.mjs';
10
11 const RULES_BLOB_GLOBAL = '__knowtation_gateway_ingest_rules_blob';
12 const IDEM_BLOB_GLOBAL = '__knowtation_gateway_ingest_idempotency_blob';
13 const RULES_BLOB_KEY = 'automation-ingest-rules-v1';
14 const IDEM_BLOB_KEY = 'automation-ingest-idempotency-v1';
15 export const INGEST_STORE_UNAVAILABLE = 'AGENT_CREDENTIAL_STORE_UNAVAILABLE';
16
17 let projectRoot;
18 try {
19 const __dirname = path.dirname(fileURLToPath(import.meta.url));
20 projectRoot = path.resolve(__dirname, '..', '..');
21 } catch (_) {
22 projectRoot = process.cwd();
23 }
24
25 function isNetlify() {
26 return typeof process.env.NETLIFY === 'string' && process.env.NETLIFY.length > 0;
27 }
28
29 function resolveDataDir(override) {
30 if (override) return override;
31 return process.env.KNOWTATION_GATEWAY_DATA_DIR || path.join(projectRoot, 'data');
32 }
33
34 function rulesFilePath(override) {
35 return path.join(resolveDataDir(override), 'automation_ingest_rules.json');
36 }
37
38 function idempotencyFilePath(override) {
39 return path.join(resolveDataDir(override), 'automation_ingest_idempotency.json');
40 }
41
42 function wrapStoreError(e) {
43 if (e && e.code === INGEST_STORE_UNAVAILABLE) return e;
44 const err = new Error(e && e.message ? e.message : 'automation ingest store I/O failed');
45 err.code = INGEST_STORE_UNAVAILABLE;
46 err.status = 503;
47 return err;
48 }
49
50 function emptyRulesEnvelope() {
51 return { version: 1, subs: {} };
52 }
53
54 function emptyIdempotencyEnvelope() {
55 return { version: 1, entries: {} };
56 }
57
58 function resolveRulesBackend() {
59 if (isNetlify()) {
60 const store = globalThis[RULES_BLOB_GLOBAL];
61 if (!store) {
62 const err = new Error('ingest rules blob global missing on Netlify');
63 err.code = INGEST_STORE_UNAVAILABLE;
64 err.status = 503;
65 throw err;
66 }
67 return { kind: 'blob', store };
68 }
69 const store = globalThis[RULES_BLOB_GLOBAL];
70 if (store) return { kind: 'blob', store };
71 return { kind: 'file' };
72 }
73
74 function resolveIdemBackend() {
75 if (isNetlify()) {
76 const store = globalThis[IDEM_BLOB_GLOBAL];
77 if (!store) {
78 const err = new Error('ingest idempotency blob global missing on Netlify');
79 err.code = INGEST_STORE_UNAVAILABLE;
80 err.status = 503;
81 throw err;
82 }
83 return { kind: 'blob', store };
84 }
85 const store = globalThis[IDEM_BLOB_GLOBAL];
86 if (store) return { kind: 'blob', store };
87 return { kind: 'file' };
88 }
89
90 async function readJsonBlob(store, key) {
91 if (!store || typeof store.get !== 'function') return null;
92 const raw = await store.get(key);
93 if (raw == null) return null;
94 if (typeof raw === 'string') {
95 if (!raw.trim()) return null;
96 return JSON.parse(raw);
97 }
98 if (typeof raw === 'object') return raw;
99 return null;
100 }
101
102 async function writeJsonBlob(store, key, value) {
103 if (!store || typeof store.set !== 'function') {
104 const err = new Error('ingest blob set unavailable');
105 err.code = INGEST_STORE_UNAVAILABLE;
106 err.status = 503;
107 throw err;
108 }
109 await store.set(key, JSON.stringify(value));
110 }
111
112 async function readJsonFile(fp) {
113 try {
114 const raw = await fs.readFile(fp, 'utf8');
115 if (!raw.trim()) return null;
116 return JSON.parse(raw);
117 } catch (e) {
118 if (e && e.code === 'ENOENT') return null;
119 throw e;
120 }
121 }
122
123 async function writeJsonFile(fp, value) {
124 await fs.mkdir(path.dirname(fp), { recursive: true });
125 await fs.writeFile(fp, JSON.stringify(value, null, 2), 'utf8');
126 }
127
128 async function loadRulesEnvelope(dataDir) {
129 try {
130 const backend = resolveRulesBackend();
131 const raw =
132 backend.kind === 'blob'
133 ? await readJsonBlob(backend.store, RULES_BLOB_KEY)
134 : await readJsonFile(rulesFilePath(dataDir));
135 if (!raw || typeof raw !== 'object') return emptyRulesEnvelope();
136 const subs = raw.subs && typeof raw.subs === 'object' ? raw.subs : {};
137 return { version: 1, subs };
138 } catch (e) {
139 throw wrapStoreError(e);
140 }
141 }
142
143 async function saveRulesEnvelope(envelope, dataDir) {
144 try {
145 const backend = resolveRulesBackend();
146 if (backend.kind === 'blob') {
147 await writeJsonBlob(backend.store, RULES_BLOB_KEY, envelope);
148 return;
149 }
150 await writeJsonFile(rulesFilePath(dataDir), envelope);
151 } catch (e) {
152 throw wrapStoreError(e);
153 }
154 }
155
156 async function loadIdempotencyEnvelope(dataDir) {
157 try {
158 const backend = resolveIdemBackend();
159 const raw =
160 backend.kind === 'blob'
161 ? await readJsonBlob(backend.store, IDEM_BLOB_KEY)
162 : await readJsonFile(idempotencyFilePath(dataDir));
163 if (!raw || typeof raw !== 'object') return emptyIdempotencyEnvelope();
164 const entries = raw.entries && typeof raw.entries === 'object' ? raw.entries : {};
165 return { version: 1, entries };
166 } catch (e) {
167 throw wrapStoreError(e);
168 }
169 }
170
171 async function saveIdempotencyEnvelope(envelope, dataDir) {
172 try {
173 const backend = resolveIdemBackend();
174 if (backend.kind === 'blob') {
175 await writeJsonBlob(backend.store, IDEM_BLOB_KEY, envelope);
176 return;
177 }
178 await writeJsonFile(idempotencyFilePath(dataDir), envelope);
179 } catch (e) {
180 throw wrapStoreError(e);
181 }
182 }
183
184 /**
185 * @param {string} sub
186 * @returns {Promise<{ rules: object[], templates: object[] }>}
187 */
188 export async function loadIngestRulesForSub(sub, dataDir) {
189 const sid = String(sub || '');
190 const envelope = await loadRulesEnvelope(dataDir);
191 const row = envelope.subs[sid];
192 const rules = row && Array.isArray(row.rules) ? row.rules : [];
193 return { rules, templates: listPackTemplates() };
194 }
195
196 /**
197 * @param {string} sub
198 * @param {object[]} rules
199 * @param {string} [dataDir]
200 */
201 export async function saveIngestRulesForSub(sub, rules, dataDir) {
202 const sid = String(sub || '');
203 const envelope = await loadRulesEnvelope(dataDir);
204 envelope.subs[sid] = { rules: Array.isArray(rules) ? rules : [], updated_at: Date.now() };
205 await saveRulesEnvelope(envelope, dataDir);
206 return envelope.subs[sid];
207 }
208
209 /**
210 * @param {string} storeKey
211 * @param {string} [dataDir]
212 */
213 export async function getIngestIdempotency(storeKey, dataDir) {
214 const envelope = await loadIdempotencyEnvelope(dataDir);
215 const row = envelope.entries[storeKey];
216 if (!row || typeof row !== 'object') return null;
217 if (Number(row.expires_at) <= Date.now()) return null;
218 return row;
219 }
220
221 /**
222 * @param {string} storeKey
223 * @param {object} entry
224 * @param {string} [dataDir]
225 */
226 export async function putIngestIdempotency(storeKey, entry, dataDir) {
227 const envelope = await loadIdempotencyEnvelope(dataDir);
228 envelope.entries[storeKey] = entry;
229 const now = Date.now();
230 for (const [k, v] of Object.entries(envelope.entries)) {
231 if (!v || Number(v.expires_at) <= now) delete envelope.entries[k];
232 }
233 await saveIdempotencyEnvelope(envelope, dataDir);
234 return envelope.entries[storeKey];
235 }
File History 1 commit
sha256:700fafdd1afa490919f9515d660ca6e75456bcd5bb67513abcd8757a634c01f6 docs: record AIP-b SD-21 land (KN #308) Human 9 days ago