Skip to content
Closed
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
63 changes: 15 additions & 48 deletions src/lib/covenant/integrations.ts
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,3 @@
import { createClient } from "redis";

// Reuse client if it exists globally to avoid reconnecting on every request
declare global {
var _redisClient: ReturnType<typeof createClient> | undefined;
}

async function getRedisClient() {
if (!global._redisClient) {
global._redisClient = createClient({ url: process.env.REDIS_URL || "redis://localhost:6379" });
global._redisClient.on("error", (err) => console.error("Redis error:", err));
await global._redisClient.connect().catch(() => {});
}
return global._redisClient;
}

export class IntegrationUnavailable extends Error {
readonly code = "INTEGRATION_UNAVAILABLE";
}
Expand All@@ -28,11 +12,13 @@ export function requireIntegration(name: string, value: string | undefined): str
return value.replace(/\/$/, "");
}

export async function postIntegration(url: string, body: unknown, headers: Record<string, string> = {}): Promise<Record<string, unknown>> {
export async function postIntegration(
url: string,
body: unknown,
headers: Record<string, string> = {},
): Promise<Record<string, unknown>> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 3000); // 3 second fail-fast

const cacheKey = `cAPI:integration:${Buffer.from(url).toString('base64')}:${Buffer.from(JSON.stringify(body)).toString('base64')}`;
const timeoutId = setTimeout(() => controller.abort(), 3000);

try {
const response = await fetch(url, {
Expand All@@ -41,48 +27,29 @@ export async function postIntegration(url: string, body: unknown, headers: Recor
body: JSON.stringify(body),
signal: controller.signal,
});
clearTimeout(timeoutId);

if (!response.ok) {
if (response.status === 401 || response.status === 403) {
throw new AuthorityDenied(`Authority denied: HTTP ${response.status}`);
}
throw new Error(`HTTP ${response.status}`);
throw new IntegrationUnavailable(`Integration failed: HTTP ${response.status}`);
}

const result: unknown = await response.json();
if (!result || typeof result !== "object" || Array.isArray(result)) {
throw new Error("Invalid response");
throw new IntegrationUnavailable("Integration failed: invalid response");
}

// Cache successful response asynchronously
getRedisClient().then(client => {
if (client.isOpen) client.setEx(cacheKey, 3600, JSON.stringify(result)).catch(console.error);
}).catch(console.error);

return result as Record<string, unknown>;
} catch (error) {
clearTimeout(timeoutId);

if (error instanceof AuthorityDenied) {
if (error instanceof AuthorityDenied || error instanceof IntegrationUnavailable) {
throw error;
}

// Attempt to retrieve stale data
try {
const client = await getRedisClient();
if (client.isOpen) {
const cached = await client.get(cacheKey);
if (cached) {
const parsed = JSON.parse(cached);
parsed._stale = true; // Mark as stale
return parsed;
}
}
} catch (redisError) {
console.error("Failed to retrieve stale cache:", redisError);
}

throw new IntegrationUnavailable(`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`);
throw new IntegrationUnavailable(
`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`,
);
} finally {
clearTimeout(timeoutId);
}
}
121 changes: 121 additions & 0 deletions tests/integrations.fail-closed.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
AuthorityDenied,
IntegrationUnavailable,
postIntegration,
} from "../src/lib/covenant/integrations";

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

describe("postIntegration fail-closed authority boundary", () => {
it("does not replay a prior successful response after the authority service fails", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockResolvedValueOnce(new Response("upstream failure", { status: 503 }));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after a network rejection", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockRejectedValueOnce(new Error("network unavailable"));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after the three-second timeout", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockImplementationOnce((_url: string, init?: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
init?.signal?.addEventListener("abort", () => {
reject(new DOMException("The operation was aborted", "AbortError"));
});
}),
);

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-timeout" }),
).resolves.toEqual({ decision: "APPROVED" });

vi.useFakeTimers();
const timedOutRequest = postIntegration("https://cappo.example.test/authorize", {
run_id: "run-timeout",
});
const timeoutExpectation = expect(timedOutRequest).rejects.toBeInstanceOf(
IntegrationUnavailable,
);

await vi.advanceTimersByTimeAsync(3000);
await timeoutExpectation;
});

it("preserves explicit authority denial", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(new Response("denied", { status: 403 })),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-2" }),
).rejects.toBeInstanceOf(AuthorityDenied);
});

it("fails closed on invalid success payloads", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(
new Response(JSON.stringify(["not", "an", "authority", "object"]), {
status: 200,
headers: { "content-type": "application/json" },
}),
),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-3" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});
});
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
fix(governance): fail closed on integration authority outages by reprewindai-dev · Pull Request #44 · reprewindai-dev/cAPI · GitHub
Skip to content
Closed
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
63 changes: 15 additions & 48 deletions src/lib/covenant/integrations.ts
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,3 @@
import { createClient } from "redis";

// Reuse client if it exists globally to avoid reconnecting on every request
declare global {
var _redisClient: ReturnType<typeof createClient> | undefined;
}

async function getRedisClient() {
if (!global._redisClient) {
global._redisClient = createClient({ url: process.env.REDIS_URL || "redis://localhost:6379" });
global._redisClient.on("error", (err) => console.error("Redis error:", err));
await global._redisClient.connect().catch(() => {});
}
return global._redisClient;
}

export class IntegrationUnavailable extends Error {
readonly code = "INTEGRATION_UNAVAILABLE";
}
Expand All@@ -28,11 +12,13 @@ export function requireIntegration(name: string, value: string | undefined): str
return value.replace(/\/$/, "");
}

export async function postIntegration(url: string, body: unknown, headers: Record<string, string> = {}): Promise<Record<string, unknown>> {
export async function postIntegration(
url: string,
body: unknown,
headers: Record<string, string> = {},
): Promise<Record<string, unknown>> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 3000); // 3 second fail-fast

const cacheKey = `cAPI:integration:${Buffer.from(url).toString('base64')}:${Buffer.from(JSON.stringify(body)).toString('base64')}`;
const timeoutId = setTimeout(() => controller.abort(), 3000);

try {
const response = await fetch(url, {
Expand All@@ -41,48 +27,29 @@ export async function postIntegration(url: string, body: unknown, headers: Recor
body: JSON.stringify(body),
signal: controller.signal,
});
clearTimeout(timeoutId);

if (!response.ok) {
if (response.status === 401 || response.status === 403) {
throw new AuthorityDenied(`Authority denied: HTTP ${response.status}`);
}
throw new Error(`HTTP ${response.status}`);
throw new IntegrationUnavailable(`Integration failed: HTTP ${response.status}`);
}

const result: unknown = await response.json();
if (!result || typeof result !== "object" || Array.isArray(result)) {
throw new Error("Invalid response");
throw new IntegrationUnavailable("Integration failed: invalid response");
}

// Cache successful response asynchronously
getRedisClient().then(client => {
if (client.isOpen) client.setEx(cacheKey, 3600, JSON.stringify(result)).catch(console.error);
}).catch(console.error);

return result as Record<string, unknown>;
} catch (error) {
clearTimeout(timeoutId);

if (error instanceof AuthorityDenied) {
if (error instanceof AuthorityDenied || error instanceof IntegrationUnavailable) {
throw error;
}

// Attempt to retrieve stale data
try {
const client = await getRedisClient();
if (client.isOpen) {
const cached = await client.get(cacheKey);
if (cached) {
const parsed = JSON.parse(cached);
parsed._stale = true; // Mark as stale
return parsed;
}
}
} catch (redisError) {
console.error("Failed to retrieve stale cache:", redisError);
}

throw new IntegrationUnavailable(`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`);
throw new IntegrationUnavailable(
`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`,
);
} finally {
clearTimeout(timeoutId);
}
}
121 changes: 121 additions & 0 deletions tests/integrations.fail-closed.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
AuthorityDenied,
IntegrationUnavailable,
postIntegration,
} from "../src/lib/covenant/integrations";

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

describe("postIntegration fail-closed authority boundary", () => {
it("does not replay a prior successful response after the authority service fails", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockResolvedValueOnce(new Response("upstream failure", { status: 503 }));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after a network rejection", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockRejectedValueOnce(new Error("network unavailable"));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after the three-second timeout", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockImplementationOnce((_url: string, init?: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
init?.signal?.addEventListener("abort", () => {
reject(new DOMException("The operation was aborted", "AbortError"));
});
}),
);

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-timeout" }),
).resolves.toEqual({ decision: "APPROVED" });

vi.useFakeTimers();
const timedOutRequest = postIntegration("https://cappo.example.test/authorize", {
run_id: "run-timeout",
});
const timeoutExpectation = expect(timedOutRequest).rejects.toBeInstanceOf(
IntegrationUnavailable,
);

await vi.advanceTimersByTimeAsync(3000);
await timeoutExpectation;
});

it("preserves explicit authority denial", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(new Response("denied", { status: 403 })),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-2" }),
).rejects.toBeInstanceOf(AuthorityDenied);
});

it("fails closed on invalid success payloads", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(
new Response(JSON.stringify(["not", "an", "authority", "object"]), {
status: 200,
headers: { "content-type": "application/json" },
}),
),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-3" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});
});
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' fix(governance): fail closed on integration authority outages by reprewindai-dev · Pull Request #44 · reprewindai-dev/cAPI · GitHub
Skip to content
Closed
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
63 changes: 15 additions & 48 deletions src/lib/covenant/integrations.ts
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,3 @@
import { createClient } from "redis";

// Reuse client if it exists globally to avoid reconnecting on every request
declare global {
var _redisClient: ReturnType<typeof createClient> | undefined;
}

async function getRedisClient() {
if (!global._redisClient) {
global._redisClient = createClient({ url: process.env.REDIS_URL || "redis://localhost:6379" });
global._redisClient.on("error", (err) => console.error("Redis error:", err));
await global._redisClient.connect().catch(() => {});
}
return global._redisClient;
}

export class IntegrationUnavailable extends Error {
readonly code = "INTEGRATION_UNAVAILABLE";
}
Expand All@@ -28,11 +12,13 @@ export function requireIntegration(name: string, value: string | undefined): str
return value.replace(/\/$/, "");
}

export async function postIntegration(url: string, body: unknown, headers: Record<string, string> = {}): Promise<Record<string, unknown>> {
export async function postIntegration(
url: string,
body: unknown,
headers: Record<string, string> = {},
): Promise<Record<string, unknown>> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 3000); // 3 second fail-fast

const cacheKey = `cAPI:integration:${Buffer.from(url).toString('base64')}:${Buffer.from(JSON.stringify(body)).toString('base64')}`;
const timeoutId = setTimeout(() => controller.abort(), 3000);

try {
const response = await fetch(url, {
Expand All@@ -41,48 +27,29 @@ export async function postIntegration(url: string, body: unknown, headers: Recor
body: JSON.stringify(body),
signal: controller.signal,
});
clearTimeout(timeoutId);

if (!response.ok) {
if (response.status === 401 || response.status === 403) {
throw new AuthorityDenied(`Authority denied: HTTP ${response.status}`);
}
throw new Error(`HTTP ${response.status}`);
throw new IntegrationUnavailable(`Integration failed: HTTP ${response.status}`);
}

const result: unknown = await response.json();
if (!result || typeof result !== "object" || Array.isArray(result)) {
throw new Error("Invalid response");
throw new IntegrationUnavailable("Integration failed: invalid response");
}

// Cache successful response asynchronously
getRedisClient().then(client => {
if (client.isOpen) client.setEx(cacheKey, 3600, JSON.stringify(result)).catch(console.error);
}).catch(console.error);

return result as Record<string, unknown>;
} catch (error) {
clearTimeout(timeoutId);

if (error instanceof AuthorityDenied) {
if (error instanceof AuthorityDenied || error instanceof IntegrationUnavailable) {
throw error;
}

// Attempt to retrieve stale data
try {
const client = await getRedisClient();
if (client.isOpen) {
const cached = await client.get(cacheKey);
if (cached) {
const parsed = JSON.parse(cached);
parsed._stale = true; // Mark as stale
return parsed;
}
}
} catch (redisError) {
console.error("Failed to retrieve stale cache:", redisError);
}

throw new IntegrationUnavailable(`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`);
throw new IntegrationUnavailable(
`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`,
);
} finally {
clearTimeout(timeoutId);
}
}
121 changes: 121 additions & 0 deletions tests/integrations.fail-closed.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
AuthorityDenied,
IntegrationUnavailable,
postIntegration,
} from "../src/lib/covenant/integrations";

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

describe("postIntegration fail-closed authority boundary", () => {
it("does not replay a prior successful response after the authority service fails", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockResolvedValueOnce(new Response("upstream failure", { status: 503 }));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after a network rejection", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockRejectedValueOnce(new Error("network unavailable"));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after the three-second timeout", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockImplementationOnce((_url: string, init?: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
init?.signal?.addEventListener("abort", () => {
reject(new DOMException("The operation was aborted", "AbortError"));
});
}),
);

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-timeout" }),
).resolves.toEqual({ decision: "APPROVED" });

vi.useFakeTimers();
const timedOutRequest = postIntegration("https://cappo.example.test/authorize", {
run_id: "run-timeout",
});
const timeoutExpectation = expect(timedOutRequest).rejects.toBeInstanceOf(
IntegrationUnavailable,
);

await vi.advanceTimersByTimeAsync(3000);
await timeoutExpectation;
});

it("preserves explicit authority denial", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(new Response("denied", { status: 403 })),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-2" }),
).rejects.toBeInstanceOf(AuthorityDenied);
});

it("fails closed on invalid success payloads", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(
new Response(JSON.stringify(["not", "an", "authority", "object"]), {
status: 200,
headers: { "content-type": "application/json" },
}),
),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-3" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});
});
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' fix(governance): fail closed on integration authority outages by reprewindai-dev · Pull Request #44 · reprewindai-dev/cAPI · GitHub
Skip to content
Closed
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
63 changes: 15 additions & 48 deletions src/lib/covenant/integrations.ts
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,3 @@
import { createClient } from "redis";

// Reuse client if it exists globally to avoid reconnecting on every request
declare global {
var _redisClient: ReturnType<typeof createClient> | undefined;
}

async function getRedisClient() {
if (!global._redisClient) {
global._redisClient = createClient({ url: process.env.REDIS_URL || "redis://localhost:6379" });
global._redisClient.on("error", (err) => console.error("Redis error:", err));
await global._redisClient.connect().catch(() => {});
}
return global._redisClient;
}

export class IntegrationUnavailable extends Error {
readonly code = "INTEGRATION_UNAVAILABLE";
}
Expand All@@ -28,11 +12,13 @@ export function requireIntegration(name: string, value: string | undefined): str
return value.replace(/\/$/, "");
}

export async function postIntegration(url: string, body: unknown, headers: Record<string, string> = {}): Promise<Record<string, unknown>> {
export async function postIntegration(
url: string,
body: unknown,
headers: Record<string, string> = {},
): Promise<Record<string, unknown>> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 3000); // 3 second fail-fast

const cacheKey = `cAPI:integration:${Buffer.from(url).toString('base64')}:${Buffer.from(JSON.stringify(body)).toString('base64')}`;
const timeoutId = setTimeout(() => controller.abort(), 3000);

try {
const response = await fetch(url, {
Expand All@@ -41,48 +27,29 @@ export async function postIntegration(url: string, body: unknown, headers: Recor
body: JSON.stringify(body),
signal: controller.signal,
});
clearTimeout(timeoutId);

if (!response.ok) {
if (response.status === 401 || response.status === 403) {
throw new AuthorityDenied(`Authority denied: HTTP ${response.status}`);
}
throw new Error(`HTTP ${response.status}`);
throw new IntegrationUnavailable(`Integration failed: HTTP ${response.status}`);
}

const result: unknown = await response.json();
if (!result || typeof result !== "object" || Array.isArray(result)) {
throw new Error("Invalid response");
throw new IntegrationUnavailable("Integration failed: invalid response");
}

// Cache successful response asynchronously
getRedisClient().then(client => {
if (client.isOpen) client.setEx(cacheKey, 3600, JSON.stringify(result)).catch(console.error);
}).catch(console.error);

return result as Record<string, unknown>;
} catch (error) {
clearTimeout(timeoutId);

if (error instanceof AuthorityDenied) {
if (error instanceof AuthorityDenied || error instanceof IntegrationUnavailable) {
throw error;
}

// Attempt to retrieve stale data
try {
const client = await getRedisClient();
if (client.isOpen) {
const cached = await client.get(cacheKey);
if (cached) {
const parsed = JSON.parse(cached);
parsed._stale = true; // Mark as stale
return parsed;
}
}
} catch (redisError) {
console.error("Failed to retrieve stale cache:", redisError);
}

throw new IntegrationUnavailable(`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`);
throw new IntegrationUnavailable(
`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`,
);
} finally {
clearTimeout(timeoutId);
}
}
121 changes: 121 additions & 0 deletions tests/integrations.fail-closed.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
AuthorityDenied,
IntegrationUnavailable,
postIntegration,
} from "../src/lib/covenant/integrations";

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

describe("postIntegration fail-closed authority boundary", () => {
it("does not replay a prior successful response after the authority service fails", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockResolvedValueOnce(new Response("upstream failure", { status: 503 }));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after a network rejection", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockRejectedValueOnce(new Error("network unavailable"));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after the three-second timeout", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockImplementationOnce((_url: string, init?: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
init?.signal?.addEventListener("abort", () => {
reject(new DOMException("The operation was aborted", "AbortError"));
});
}),
);

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-timeout" }),
).resolves.toEqual({ decision: "APPROVED" });

vi.useFakeTimers();
const timedOutRequest = postIntegration("https://cappo.example.test/authorize", {
run_id: "run-timeout",
});
const timeoutExpectation = expect(timedOutRequest).rejects.toBeInstanceOf(
IntegrationUnavailable,
);

await vi.advanceTimersByTimeAsync(3000);
await timeoutExpectation;
});

it("preserves explicit authority denial", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(new Response("denied", { status: 403 })),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-2" }),
).rejects.toBeInstanceOf(AuthorityDenied);
});

it("fails closed on invalid success payloads", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(
new Response(JSON.stringify(["not", "an", "authority", "object"]), {
status: 200,
headers: { "content-type": "application/json" },
}),
),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-3" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});
});
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' fix(governance): fail closed on integration authority outages by reprewindai-dev · Pull Request #44 · reprewindai-dev/cAPI · GitHub
Skip to content
Closed
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
63 changes: 15 additions & 48 deletions src/lib/covenant/integrations.ts
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,3 @@
import { createClient } from "redis";

// Reuse client if it exists globally to avoid reconnecting on every request
declare global {
var _redisClient: ReturnType<typeof createClient> | undefined;
}

async function getRedisClient() {
if (!global._redisClient) {
global._redisClient = createClient({ url: process.env.REDIS_URL || "redis://localhost:6379" });
global._redisClient.on("error", (err) => console.error("Redis error:", err));
await global._redisClient.connect().catch(() => {});
}
return global._redisClient;
}

export class IntegrationUnavailable extends Error {
readonly code = "INTEGRATION_UNAVAILABLE";
}
Expand All@@ -28,11 +12,13 @@ export function requireIntegration(name: string, value: string | undefined): str
return value.replace(/\/$/, "");
}

export async function postIntegration(url: string, body: unknown, headers: Record<string, string> = {}): Promise<Record<string, unknown>> {
export async function postIntegration(
url: string,
body: unknown,
headers: Record<string, string> = {},
): Promise<Record<string, unknown>> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 3000); // 3 second fail-fast

const cacheKey = `cAPI:integration:${Buffer.from(url).toString('base64')}:${Buffer.from(JSON.stringify(body)).toString('base64')}`;
const timeoutId = setTimeout(() => controller.abort(), 3000);

try {
const response = await fetch(url, {
Expand All@@ -41,48 +27,29 @@ export async function postIntegration(url: string, body: unknown, headers: Recor
body: JSON.stringify(body),
signal: controller.signal,
});
clearTimeout(timeoutId);

if (!response.ok) {
if (response.status === 401 || response.status === 403) {
throw new AuthorityDenied(`Authority denied: HTTP ${response.status}`);
}
throw new Error(`HTTP ${response.status}`);
throw new IntegrationUnavailable(`Integration failed: HTTP ${response.status}`);
}

const result: unknown = await response.json();
if (!result || typeof result !== "object" || Array.isArray(result)) {
throw new Error("Invalid response");
throw new IntegrationUnavailable("Integration failed: invalid response");
}

// Cache successful response asynchronously
getRedisClient().then(client => {
if (client.isOpen) client.setEx(cacheKey, 3600, JSON.stringify(result)).catch(console.error);
}).catch(console.error);

return result as Record<string, unknown>;
} catch (error) {
clearTimeout(timeoutId);

if (error instanceof AuthorityDenied) {
if (error instanceof AuthorityDenied || error instanceof IntegrationUnavailable) {
throw error;
}

// Attempt to retrieve stale data
try {
const client = await getRedisClient();
if (client.isOpen) {
const cached = await client.get(cacheKey);
if (cached) {
const parsed = JSON.parse(cached);
parsed._stale = true; // Mark as stale
return parsed;
}
}
} catch (redisError) {
console.error("Failed to retrieve stale cache:", redisError);
}

throw new IntegrationUnavailable(`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`);
throw new IntegrationUnavailable(
`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`,
);
} finally {
clearTimeout(timeoutId);
}
}
121 changes: 121 additions & 0 deletions tests/integrations.fail-closed.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
AuthorityDenied,
IntegrationUnavailable,
postIntegration,
} from "../src/lib/covenant/integrations";

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

describe("postIntegration fail-closed authority boundary", () => {
it("does not replay a prior successful response after the authority service fails", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockResolvedValueOnce(new Response("upstream failure", { status: 503 }));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after a network rejection", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockRejectedValueOnce(new Error("network unavailable"));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after the three-second timeout", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockImplementationOnce((_url: string, init?: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
init?.signal?.addEventListener("abort", () => {
reject(new DOMException("The operation was aborted", "AbortError"));
});
}),
);

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-timeout" }),
).resolves.toEqual({ decision: "APPROVED" });

vi.useFakeTimers();
const timedOutRequest = postIntegration("https://cappo.example.test/authorize", {
run_id: "run-timeout",
});
const timeoutExpectation = expect(timedOutRequest).rejects.toBeInstanceOf(
IntegrationUnavailable,
);

await vi.advanceTimersByTimeAsync(3000);
await timeoutExpectation;
});

it("preserves explicit authority denial", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(new Response("denied", { status: 403 })),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-2" }),
).rejects.toBeInstanceOf(AuthorityDenied);
});

it("fails closed on invalid success payloads", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(
new Response(JSON.stringify(["not", "an", "authority", "object"]), {
status: 200,
headers: { "content-type": "application/json" },
}),
),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-3" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});
});
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' fix(governance): fail closed on integration authority outages by reprewindai-dev · Pull Request #44 · reprewindai-dev/cAPI · GitHub
Skip to content
Closed
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
63 changes: 15 additions & 48 deletions src/lib/covenant/integrations.ts
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,3 @@
import { createClient } from "redis";

// Reuse client if it exists globally to avoid reconnecting on every request
declare global {
var _redisClient: ReturnType<typeof createClient> | undefined;
}

async function getRedisClient() {
if (!global._redisClient) {
global._redisClient = createClient({ url: process.env.REDIS_URL || "redis://localhost:6379" });
global._redisClient.on("error", (err) => console.error("Redis error:", err));
await global._redisClient.connect().catch(() => {});
}
return global._redisClient;
}

export class IntegrationUnavailable extends Error {
readonly code = "INTEGRATION_UNAVAILABLE";
}
Expand All@@ -28,11 +12,13 @@ export function requireIntegration(name: string, value: string | undefined): str
return value.replace(/\/$/, "");
}

export async function postIntegration(url: string, body: unknown, headers: Record<string, string> = {}): Promise<Record<string, unknown>> {
export async function postIntegration(
url: string,
body: unknown,
headers: Record<string, string> = {},
): Promise<Record<string, unknown>> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 3000); // 3 second fail-fast

const cacheKey = `cAPI:integration:${Buffer.from(url).toString('base64')}:${Buffer.from(JSON.stringify(body)).toString('base64')}`;
const timeoutId = setTimeout(() => controller.abort(), 3000);

try {
const response = await fetch(url, {
Expand All@@ -41,48 +27,29 @@ export async function postIntegration(url: string, body: unknown, headers: Recor
body: JSON.stringify(body),
signal: controller.signal,
});
clearTimeout(timeoutId);

if (!response.ok) {
if (response.status === 401 || response.status === 403) {
throw new AuthorityDenied(`Authority denied: HTTP ${response.status}`);
}
throw new Error(`HTTP ${response.status}`);
throw new IntegrationUnavailable(`Integration failed: HTTP ${response.status}`);
}

const result: unknown = await response.json();
if (!result || typeof result !== "object" || Array.isArray(result)) {
throw new Error("Invalid response");
throw new IntegrationUnavailable("Integration failed: invalid response");
}

// Cache successful response asynchronously
getRedisClient().then(client => {
if (client.isOpen) client.setEx(cacheKey, 3600, JSON.stringify(result)).catch(console.error);
}).catch(console.error);

return result as Record<string, unknown>;
} catch (error) {
clearTimeout(timeoutId);

if (error instanceof AuthorityDenied) {
if (error instanceof AuthorityDenied || error instanceof IntegrationUnavailable) {
throw error;
}

// Attempt to retrieve stale data
try {
const client = await getRedisClient();
if (client.isOpen) {
const cached = await client.get(cacheKey);
if (cached) {
const parsed = JSON.parse(cached);
parsed._stale = true; // Mark as stale
return parsed;
}
}
} catch (redisError) {
console.error("Failed to retrieve stale cache:", redisError);
}

throw new IntegrationUnavailable(`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`);
throw new IntegrationUnavailable(
`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`,
);
} finally {
clearTimeout(timeoutId);
}
}
121 changes: 121 additions & 0 deletions tests/integrations.fail-closed.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
AuthorityDenied,
IntegrationUnavailable,
postIntegration,
} from "../src/lib/covenant/integrations";

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

describe("postIntegration fail-closed authority boundary", () => {
it("does not replay a prior successful response after the authority service fails", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockResolvedValueOnce(new Response("upstream failure", { status: 503 }));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after a network rejection", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockRejectedValueOnce(new Error("network unavailable"));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after the three-second timeout", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockImplementationOnce((_url: string, init?: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
init?.signal?.addEventListener("abort", () => {
reject(new DOMException("The operation was aborted", "AbortError"));
});
}),
);

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-timeout" }),
).resolves.toEqual({ decision: "APPROVED" });

vi.useFakeTimers();
const timedOutRequest = postIntegration("https://cappo.example.test/authorize", {
run_id: "run-timeout",
});
const timeoutExpectation = expect(timedOutRequest).rejects.toBeInstanceOf(
IntegrationUnavailable,
);

await vi.advanceTimersByTimeAsync(3000);
await timeoutExpectation;
});

it("preserves explicit authority denial", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(new Response("denied", { status: 403 })),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-2" }),
).rejects.toBeInstanceOf(AuthorityDenied);
});

it("fails closed on invalid success payloads", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(
new Response(JSON.stringify(["not", "an", "authority", "object"]), {
status: 200,
headers: { "content-type": "application/json" },
}),
),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-3" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});
});
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); })(); fix(governance): fail closed on integration authority outages by reprewindai-dev · Pull Request #44 · reprewindai-dev/cAPI · GitHub
Skip to content
Closed
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
63 changes: 15 additions & 48 deletions src/lib/covenant/integrations.ts
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,3 @@
import { createClient } from "redis";

// Reuse client if it exists globally to avoid reconnecting on every request
declare global {
var _redisClient: ReturnType<typeof createClient> | undefined;
}

async function getRedisClient() {
if (!global._redisClient) {
global._redisClient = createClient({ url: process.env.REDIS_URL || "redis://localhost:6379" });
global._redisClient.on("error", (err) => console.error("Redis error:", err));
await global._redisClient.connect().catch(() => {});
}
return global._redisClient;
}

export class IntegrationUnavailable extends Error {
readonly code = "INTEGRATION_UNAVAILABLE";
}
Expand All@@ -28,11 +12,13 @@ export function requireIntegration(name: string, value: string | undefined): str
return value.replace(/\/$/, "");
}

export async function postIntegration(url: string, body: unknown, headers: Record<string, string> = {}): Promise<Record<string, unknown>> {
export async function postIntegration(
url: string,
body: unknown,
headers: Record<string, string> = {},
): Promise<Record<string, unknown>> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 3000); // 3 second fail-fast

const cacheKey = `cAPI:integration:${Buffer.from(url).toString('base64')}:${Buffer.from(JSON.stringify(body)).toString('base64')}`;
const timeoutId = setTimeout(() => controller.abort(), 3000);

try {
const response = await fetch(url, {
Expand All@@ -41,48 +27,29 @@ export async function postIntegration(url: string, body: unknown, headers: Recor
body: JSON.stringify(body),
signal: controller.signal,
});
clearTimeout(timeoutId);

if (!response.ok) {
if (response.status === 401 || response.status === 403) {
throw new AuthorityDenied(`Authority denied: HTTP ${response.status}`);
}
throw new Error(`HTTP ${response.status}`);
throw new IntegrationUnavailable(`Integration failed: HTTP ${response.status}`);
}

const result: unknown = await response.json();
if (!result || typeof result !== "object" || Array.isArray(result)) {
throw new Error("Invalid response");
throw new IntegrationUnavailable("Integration failed: invalid response");
}

// Cache successful response asynchronously
getRedisClient().then(client => {
if (client.isOpen) client.setEx(cacheKey, 3600, JSON.stringify(result)).catch(console.error);
}).catch(console.error);

return result as Record<string, unknown>;
} catch (error) {
clearTimeout(timeoutId);

if (error instanceof AuthorityDenied) {
if (error instanceof AuthorityDenied || error instanceof IntegrationUnavailable) {
throw error;
}

// Attempt to retrieve stale data
try {
const client = await getRedisClient();
if (client.isOpen) {
const cached = await client.get(cacheKey);
if (cached) {
const parsed = JSON.parse(cached);
parsed._stale = true; // Mark as stale
return parsed;
}
}
} catch (redisError) {
console.error("Failed to retrieve stale cache:", redisError);
}

throw new IntegrationUnavailable(`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`);
throw new IntegrationUnavailable(
`Integration failed: ${error instanceof Error ? error.message : "Unknown error"}`,
);
} finally {
clearTimeout(timeoutId);
}
}
121 changes: 121 additions & 0 deletions tests/integrations.fail-closed.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
AuthorityDenied,
IntegrationUnavailable,
postIntegration,
} from "../src/lib/covenant/integrations";

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

describe("postIntegration fail-closed authority boundary", () => {
it("does not replay a prior successful response after the authority service fails", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockResolvedValueOnce(new Response("upstream failure", { status: 503 }));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-1" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after a network rejection", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockRejectedValueOnce(new Error("network unavailable"));

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).resolves.toEqual({ decision: "APPROVED" });

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-network" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});

it("does not replay a prior successful response after the three-second timeout", async () => {
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(JSON.stringify({ decision: "APPROVED" }), {
status: 200,
headers: { "content-type": "application/json" },
}),
)
.mockImplementationOnce((_url: string, init?: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
init?.signal?.addEventListener("abort", () => {
reject(new DOMException("The operation was aborted", "AbortError"));
});
}),
);

vi.stubGlobal("fetch", fetchMock);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-timeout" }),
).resolves.toEqual({ decision: "APPROVED" });

vi.useFakeTimers();
const timedOutRequest = postIntegration("https://cappo.example.test/authorize", {
run_id: "run-timeout",
});
const timeoutExpectation = expect(timedOutRequest).rejects.toBeInstanceOf(
IntegrationUnavailable,
);

await vi.advanceTimersByTimeAsync(3000);
await timeoutExpectation;
});

it("preserves explicit authority denial", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(new Response("denied", { status: 403 })),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-2" }),
).rejects.toBeInstanceOf(AuthorityDenied);
});

it("fails closed on invalid success payloads", async () => {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue(
new Response(JSON.stringify(["not", "an", "authority", "object"]), {
status: 200,
headers: { "content-type": "application/json" },
}),
),
);

await expect(
postIntegration("https://cappo.example.test/authorize", { run_id: "run-3" }),
).rejects.toBeInstanceOf(IntegrationUnavailable);
});
});
Loading