Skip to content

Commit 241c273

Browse files
prisisclaude
andcommitted
feat(cloud): observability ingest pipeline — issues + incidents (Phase 3) (#140)
* feat(cloud): add the observability ingest pipeline (issues + incidents) Phase 3 of the observability plan — durable, cross-deployment monitoring in the Lunora Cloud control plane, fed by the Phase 2 OTLP transport. - ingest: `POST /v1/telemetry` accepts OTLP-over-HTTP/JSON from the tenant `otlpSink` and the container exporter, decodes the error spans (`src/telemetry/otlp.ts`), and folds them into grouped issues/incidents through a deploy-key-authorized `telemetry.ingest` mutation. Synchronous — the cloud app has no queue producer binding, so ingest inserts to D1 directly (like `usage.ingest`); auth reuses `authorizeDeployKey`, not the plaintext admin token. - store: `issues` + `incidents` `.global()` D1 tables, fingerprinted with `@lunora/fingerprint` (the same hash the local Studio computes, so a local Issue and a cloud Issue are one object); `lunora/{issues,incidents}.ts` member-authorized read/triage functions. A `TelemetryStore` adapter (`src/telemetry/store.ts`) owns the non-relational side — AE metrics plus a guarded Pipeline→R2 archive, each a no-op without its binding. - dashboard: hosted `IssuesSection` / `IncidentsSection`, gated behind the `logStreams` entitlement, wired into `OrganizationDashboard`. - bindings: a `TELEMETRY` AE dataset + `TELEMETRY_BUCKET` R2 bucket. Vendors `@lunora/fingerprint` (Phase 1, not yet merged) so this stacks on the cloud branch; the graft folds away once Phase 1 lands on alpha. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018sRFb1136YE8KDmDbFMYmm * feat(cloud): observability alerts — rules + delivery (Phase 4) (#141) * feat(cloud): add observability alerts — rules, firing + delivery Phase 4 of the observability plan (the "watches while you sleep" tier), stacked on the Phase 3 ingest. - schema: `alertRules` (name, target issue/incident, threshold, channel email/webhook, destination, enabled) + `alerts` (fired-alert audit trail with firing→delivered state, notification denormalized). - firing: the telemetry `ingest` mutation loads the org's enabled rules and fires each the first time a source's count crosses its threshold (`before < threshold <= after`, so exactly once), inserting a `firing` alert row. The pure crossing/render logic lives in `src/telemetry/alerts.ts` (unit-tested), mirroring how `usage.ingest` delegates to `evaluateSpendCap`. - delivery: the `/v1/telemetry` edge handler delivers fired alerts best-effort (email via `@lunora/mail`, webhook via JSON POST) then stamps them delivered — never blocking or failing ingest. - functions: `alerts.{rules,createRule,setRuleEnabled,deleteRule,list, markDelivered}` (member-authed reads/writes; deploy-key-authed markDelivered). - dashboard: `AlertsSection` (manage rules + recent fired alerts), gated behind the `logStreams` entitlement, wired into `OrganizationDashboard`. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018sRFb1136YE8KDmDbFMYmm * fix(cloud): validate webhook alert destinations against SSRF An alert rule's webhook `destination` is `fetch`ed by the control plane when the alert fires, so an owner/admin could otherwise aim it at internal infrastructure (loopback, RFC-1918, the 169.254.169.254 metadata IP, …) — server-side request forgery. Add a pure `isSafeWebhookUrl` guard (https only, public host, no embedded credentials, no loopback/private/link-local IPv4 or IPv6) enforced both at `createRule` (reject the rule) and in `deliverAlert` (never fetch an unsafe target — defense in depth for any rule created before this guard). String-level, so it can't defeat DNS rebinding, but it blocks the direct-address cases. Unit-tested. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018sRFb1136YE8KDmDbFMYmm --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> * fix(cloud): harden webhook SSRF guard Two SSRF gaps in the Observability alert delivery path: - deliverAlert followed webhook redirects, so a destination that passes isSafeWebhookUrl could 3xx-redirect to an internal address (e.g. the metadata IP). Set redirect: "manual" and reject 3xx responses. - isSafeWebhookUrl let IPv4-mapped IPv6 (::ffff:169.254.169.254, which the URL parser compresses to ::ffff:7f00:1) and the unspecified address (::) through. Reject the whole ::-prefixed non-global class. Numeric IPv4 forms (2130706433, 0x7f000001, 0177.0.0.1) were already blocked via WHATWG URL normalization; added as regression tests. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017hfLmCwH5xMfz7L73LRPFj --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 7c630e6 commit 241c273

24 files changed

Lines changed: 5759 additions & 1355 deletions
Lines changed: 198 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,198 @@
1+
import type { AnalyticsEngineDatasetLike } from "@lunora/bindings/analytics";
2+
import type { PipelineBindingLike } from "@lunora/bindings/pipelines";
3+
import { describe, expect, it } from "vitest";
4+
5+
import { crossesThreshold, isSafeWebhookUrl, renderAlert } from "../src/telemetry/alerts";
6+
import type { OtlpTracePayload } from "../src/telemetry/otlp";
7+
import { decodeTelemetryEvents } from "../src/telemetry/otlp";
8+
import { createCloudflareTelemetryStore } from "../src/telemetry/store";
9+
10+
/** Build an OTLP `KeyValue[]` from a flat string map. */
11+
const attributes = (map: Record<string, string>): { key: string; value: { stringValue: string } }[] =>
12+
Object.entries(map).map(([key, stringValue]) => {
13+
return { key, value: { stringValue } };
14+
});
15+
16+
/** One error span (status.code 2) under the given instrumentation scope. */
17+
const errorSpan = (options: { attributes?: Record<string, string>; endMs?: number; message?: string; name?: string }) => {
18+
return {
19+
attributes: attributes(options.attributes ?? {}),
20+
endTimeUnixNano: `${String(options.endMs ?? 1_700_000_000_000)}000000`,
21+
name: options.name ?? "op",
22+
status: { code: 2, message: options.message },
23+
};
24+
};
25+
26+
/** Wrap spans into an OTLP trace payload with one resource + scope. */
27+
const payload = (scopeName: string, spans: unknown[], serviceName?: string): OtlpTracePayload => {
28+
return {
29+
resourceSpans: [
30+
{
31+
resource: serviceName === undefined ? {} : { attributes: attributes({ "service.name": serviceName }) },
32+
scopeSpans: [{ scope: { name: scopeName }, spans: spans as never }],
33+
},
34+
],
35+
};
36+
};
37+
38+
describe(decodeTelemetryEvents, () => {
39+
it("decodes a worker error span into a normalized error event", () => {
40+
const events = decodeTelemetryEvents(
41+
payload("@lunora/runtime", [
42+
errorSpan({
43+
attributes: { "error.type": "CONFLICT", "lunora.function_path": "messages:send" },
44+
endMs: 1_700_000_000_000,
45+
message: "duplicate key",
46+
}),
47+
]),
48+
);
49+
50+
expect(events).toHaveLength(1);
51+
expect(events[0]).toStrictEqual({
52+
code: "CONFLICT",
53+
functionPath: "messages:send",
54+
kind: "error",
55+
message: "duplicate key",
56+
ts: 1_700_000_000_000,
57+
});
58+
});
59+
60+
it("decodes a container error span with the service name as the culprit", () => {
61+
const events = decodeTelemetryEvents(
62+
payload("@lunora/container", [errorSpan({ attributes: { "error.type": "OOMKilled" }, message: "exit 137" })], "transcoder"),
63+
);
64+
65+
expect(events).toHaveLength(1);
66+
expect(events[0]).toMatchObject({
67+
code: "OOMKilled",
68+
container: "transcoder",
69+
functionPath: "container:transcoder",
70+
kind: "container",
71+
message: "exit 137",
72+
});
73+
});
74+
75+
it("falls back to the span name when the function-path attribute is absent", () => {
76+
const events = decodeTelemetryEvents(payload("@lunora/runtime", [errorSpan({ message: "boom", name: "users:get" })]));
77+
78+
expect(events[0]?.functionPath).toBe("users:get");
79+
});
80+
81+
it("skips non-error spans and tolerates a malformed payload", () => {
82+
const okSpan = { ...errorSpan({ message: "fine" }), status: { code: 1 } };
83+
84+
expect(decodeTelemetryEvents(payload("@lunora/runtime", [okSpan]))).toHaveLength(0);
85+
expect(decodeTelemetryEvents({})).toHaveLength(0);
86+
expect(decodeTelemetryEvents({ resourceSpans: [{}] })).toHaveLength(0);
87+
});
88+
89+
it("converts nanosecond timestamps to epoch millis exactly", () => {
90+
const events = decodeTelemetryEvents(payload("@lunora/runtime", [errorSpan({ endMs: 1_699_999_999_123, message: "x" })]));
91+
92+
expect(events[0]?.ts).toBe(1_699_999_999_123);
93+
});
94+
});
95+
96+
describe(createCloudflareTelemetryStore, () => {
97+
it("no-ops without bindings", async () => {
98+
const store = createCloudflareTelemetryStore({});
99+
100+
expect(() => {
101+
store.recordCounts({ incidents: 1, issues: 2, organizationId: "org_1" });
102+
}).not.toThrow();
103+
await expect(store.archiveEvents([{ functionPath: "a:b", kind: "error", message: "m", ts: 1 }])).resolves.toBeUndefined();
104+
});
105+
106+
it("writes one AE data point with the issue/incident counts", () => {
107+
const points: unknown[] = [];
108+
const TELEMETRY = {
109+
writeDataPoint: (event: unknown) => {
110+
points.push(event);
111+
},
112+
} as unknown as AnalyticsEngineDatasetLike;
113+
114+
createCloudflareTelemetryStore({ TELEMETRY }).recordCounts({ incidents: 3, issues: 7, organizationId: "org_42" });
115+
116+
expect(points).toStrictEqual([{ blobs: ["telemetry.ingest", "org_42"], doubles: [7, 3], indexes: ["org_42"] }]);
117+
});
118+
119+
it("archives decoded events through the pipeline binding", async () => {
120+
const batches: unknown[][] = [];
121+
const TELEMETRY_PIPELINE = {
122+
send: async (records: unknown[]) => {
123+
batches.push(records);
124+
},
125+
} as unknown as PipelineBindingLike;
126+
127+
await createCloudflareTelemetryStore({ TELEMETRY_PIPELINE }).archiveEvents([{ functionPath: "a:b", kind: "error", message: "m", ts: 1 }]);
128+
129+
expect(batches).toHaveLength(1);
130+
expect(batches[0]).toHaveLength(1);
131+
});
132+
});
133+
134+
describe(crossesThreshold, () => {
135+
it("fires only on the ingest that first reaches the threshold", () => {
136+
expect(crossesThreshold(3, 5, 5)).toBe(true); // 3 → 5 crosses 5
137+
expect(crossesThreshold(5, 7, 5)).toBe(false); // already over — fired earlier
138+
expect(crossesThreshold(0, 4, 5)).toBe(false); // not there yet
139+
expect(crossesThreshold(0, 5, 5)).toBe(true); // brand-new source straight to threshold
140+
});
141+
});
142+
143+
describe(renderAlert, () => {
144+
it("renders an issue subject + body from the rule and source", () => {
145+
const rendered = renderAlert(
146+
{ name: "High error rate", target: "issue" },
147+
{ count: 12, culprit: "messages:send", sampleMessage: "duplicate key", title: "duplicate key" },
148+
);
149+
150+
expect(rendered.subject).toBe("[Lunora] High error rate: duplicate key");
151+
expect(rendered.body).toContain('Issue "duplicate key" (messages:send) reached 12 events');
152+
expect(rendered.body).toContain("Sample: duplicate key");
153+
});
154+
155+
it("labels an incident source as an Incident", () => {
156+
const rendered = renderAlert(
157+
{ name: "Crashes", target: "incident" },
158+
{ count: 3, culprit: "container:transcoder", sampleMessage: "exit 137", title: "exit 137" },
159+
);
160+
161+
expect(rendered.body).toContain('Incident "exit 137" (container:transcoder)');
162+
});
163+
});
164+
165+
describe(isSafeWebhookUrl, () => {
166+
it("accepts an https URL to a public host", () => {
167+
expect(isSafeWebhookUrl("https://hooks.example.com/lunora")).toBe(true);
168+
expect(isSafeWebhookUrl("https://203.0.113.10/hook")).toBe(true);
169+
expect(isSafeWebhookUrl("https://[2606:4700::1111]/hook")).toBe(true); // public IPv6
170+
});
171+
172+
it("rejects SSRF-prone destinations", () => {
173+
for (const bad of [
174+
"http://hooks.example.com/x", // not https
175+
"https://someone@hooks.example.com/x", // embedded credentials (userinfo)
176+
"https://localhost/x",
177+
"https://svc.internal/x",
178+
"https://api.local/x",
179+
"https://127.0.0.1/x", // loopback
180+
"https://169.254.169.254/latest/meta-data", // cloud metadata
181+
"https://10.0.0.5/x",
182+
"https://192.168.1.1/x",
183+
"https://172.16.0.9/x",
184+
"https://100.64.0.1/x", // CGNAT
185+
"https://[::1]/x", // IPv6 loopback
186+
"https://[::]/x", // IPv6 unspecified
187+
"https://[::ffff:169.254.169.254]/x", // IPv4-mapped IPv6 → metadata IP
188+
"https://[::ffff:127.0.0.1]/x", // IPv4-mapped IPv6 → loopback
189+
"https://2130706433/x", // decimal-encoded 127.0.0.1
190+
"https://0x7f000001/x", // hex-encoded 127.0.0.1
191+
"https://0177.0.0.1/x", // octal-encoded loopback
192+
"ftp://example.com/x",
193+
"not a url",
194+
]) {
195+
expect(isSafeWebhookUrl(bad)).toBe(false);
196+
}
197+
});
198+
});

‎apps/cloud/lunora/_generated/api.ts‎

Lines changed: 22 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,14 @@ import type { FunctionReference } from "@lunora/client";
77
import type { Id } from "./dataModel.js";
88

99
export interface ApiTypes {
10+
alerts: {
11+
createRule: FunctionReference<"mutation", { channel: "email" | "webhook"; destination: string; name: string; organizationId: Id<"organizations">; target: "issue" | "incident"; threshold: number }, Id<"alertRules">>;
12+
deleteRule: FunctionReference<"mutation", { id: Id<"alertRules">; organizationId: Id<"organizations"> }, Id<"alertRules">>;
13+
list: FunctionReference<"query", { organizationId: Id<"organizations"> }, { _id: Id<"alerts">; channel: "email" | "webhook"; createdAt: number; deliveredAt?: number; destination: string; status: "delivered" | "failed" | "firing"; subject: string; target: "incident" | "issue" }[]>;
14+
markDelivered: FunctionReference<"mutation", { deployKey: string; ids: Array<Id<"alerts">>; organizationId: Id<"organizations"> }, { delivered: number; }>;
15+
rules: FunctionReference<"query", { organizationId: Id<"organizations"> }, { _id: Id<"alertRules">; channel: "email" | "webhook"; createdAt: number; destination: string; enabled: boolean; name: string; organizationId: Id<"organizations">; target: "incident" | "issue"; threshold: number }[]>;
16+
setRuleEnabled: FunctionReference<"mutation", { enabled: boolean; id: Id<"alertRules">; organizationId: Id<"organizations"> }, Id<"alertRules">>;
17+
};
1018
audit_log: {
1119
list: FunctionReference<"query", { organizationId: Id<"organizations"> }, { _id: Id<"auditLog">; action: string; actorUserId: string; createdAt: number; organizationId: Id<"organizations">; target?: string }[]>;
1220
record: FunctionReference<"mutation", { action: string; organizationId: Id<"organizations">; target?: string }, Id<"auditLog">>;
@@ -19,7 +27,7 @@ export interface ApiTypes {
1927
subscription: FunctionReference<"query", { organizationId: Id<"organizations"> }, { cancelAtPeriodEnd?: false | true; currentPeriodEnd?: number; priceId: string; provider: string; referenceId: string; state: string }[]>;
2028
};
2129
builds: {
22-
listByProject: FunctionReference<"query", { organizationId: Id<"organizations">; projectId: Id<"projects"> }, { _id: Id<"builds">; branch: string; bundleHash?: string; commitSha: string; createdAt: number; organizationId: Id<"organizations">; processingBy?: string; processingStartedAt?: number; projectId: Id<"projects">; status: "building" | "failed" | "pending" | "successful" }[]>;
30+
listByProject: FunctionReference<"query", { organizationId: Id<"organizations">; projectId: Id<"projects"> }, { _id: Id<"builds">; branch: string; bundleHash?: string; commitSha: string; createdAt: number; organizationId: Id<"organizations">; processingBy?: string; processingStartedAt?: number; projectId: Id<"projects">; status: "failed" | "building" | "pending" | "successful" }[]>;
2331
logs: FunctionReference<"query", { afterCreatedAt?: number; buildId: Id<"builds">; organizationId: Id<"organizations"> }, { createdAt: number; level: "error" | "info"; line: string; }[]>;
2432
recordPush: FunctionReference<"mutation", { branch: string; commitSha: string; installationId: number; repository: string }, { buildId: Id<"builds">; reused: boolean; } | null>;
2533
};
@@ -37,7 +45,7 @@ export interface ApiTypes {
3745
activate: FunctionReference<"mutation", { deployKey?: string; id: Id<"deployments"> }, void>;
3846
adminTarget: FunctionReference<"query", { deploymentId: Id<"deployments">; organizationId: Id<"organizations"> }, { adminToken: string; url: string; } | null>;
3947
create: FunctionReference<"mutation", { adminToken?: string; branch?: string; cronSpecs?: Array<string>; deployKey?: string; kind: "production" | "preview" | "dev"; organizationId: Id<"organizations">; projectId: Id<"projects">; runtimeVersion?: string; scriptName: string }, { deploymentId: Id<"deployments">; scriptName: string; version: number; }>;
40-
listByProject: FunctionReference<"query", { organizationId: Id<"organizations">; projectId: Id<"projects"> }, { _id: Id<"deployments">; adminToken?: string; alias?: string; branch?: string; bundleHash?: string; createdAt: number; createdBy: string; expiresAt?: number; kind: "dev" | "preview" | "production"; organizationId: Id<"organizations">; projectId: Id<"projects">; scriptName: string; status: "building" | "failed" | "destroyed" | "live" | "provisioning" | "queued" | "superseded" | "verifying"; updatedAt: number; url?: string; version?: number }[]>;
48+
listByProject: FunctionReference<"query", { organizationId: Id<"organizations">; projectId: Id<"projects"> }, { _id: Id<"deployments">; adminToken?: string; alias?: string; branch?: string; bundleHash?: string; createdAt: number; createdBy: string; expiresAt?: number; kind: "dev" | "preview" | "production"; organizationId: Id<"organizations">; projectId: Id<"projects">; scriptName: string; status: "failed" | "building" | "destroyed" | "live" | "provisioning" | "queued" | "superseded" | "verifying"; updatedAt: number; url?: string; version?: number }[]>;
4149
planForScript: FunctionReference<"query", { scriptName: string }, { plan: string; }>;
4250
rollback: FunctionReference<"mutation", { deployKey?: string; id: Id<"deployments">; organizationId: Id<"organizations"> }, { scriptName: string; version?: number; }>;
4351
routeForAlias: FunctionReference<"query", { alias: string }, { scriptName: string; } | null>;
@@ -57,12 +65,20 @@ export interface ApiTypes {
5765
record: FunctionReference<"mutation", { accountLogin: string; installationId: number }, Id<"githubInstallations">>;
5866
remove: FunctionReference<"mutation", { installationId: number }, void>;
5967
};
68+
incidents: {
69+
list: FunctionReference<"query", { organizationId: Id<"organizations"> }, { _id: Id<"incidents">; closedAt?: number; container?: string; count: number; instance?: string; kind: "crash_loop" | "error_spike" | "oom"; lastSeen: number; openedAt: number; organizationId: Id<"organizations">; status: "open" | "resolved"; title: string }[]>;
70+
setStatus: FunctionReference<"mutation", { id: Id<"incidents">; organizationId: Id<"organizations">; status: unknown }, Id<"incidents">>;
71+
};
6072
invitations: {
6173
accept: FunctionReference<"mutation", { token: string }, { organizationId: Id<"organizations">; }>;
6274
invite: FunctionReference<"mutation", { email: string; organizationId: Id<"organizations"> }, { id: Id<"invitations">; token: string; }>;
63-
list: FunctionReference<"query", { organizationId: Id<"organizations"> }, { organizationId: Id<"organizations">; status: "pending" | "accepted" | "revoked"; _id: Id<"invitations">; createdAt: number; email: string; expiresAt: number; invitedBy: string; role: "admin" | "member" | "owner" | "viewer" }[]>;
75+
list: FunctionReference<"query", { organizationId: Id<"organizations"> }, { organizationId: Id<"organizations">; email: string; status: "pending" | "accepted" | "revoked"; _id: Id<"invitations">; createdAt: number; expiresAt: number; invitedBy: string; role: "admin" | "member" | "owner" | "viewer" }[]>;
6476
revoke: FunctionReference<"mutation", { id: Id<"invitations">; organizationId: Id<"organizations"> }, void>;
6577
};
78+
issues: {
79+
list: FunctionReference<"query", { organizationId: Id<"organizations"> }, { _id: Id<"issues">; count: number; culprit: string; firstSeen: number; hash: string; lastSeen: number; organizationId: Id<"organizations">; sampleMessage: string; status: "open" | "resolved"; title: string }[]>;
80+
setStatus: FunctionReference<"mutation", { id: Id<"issues">; organizationId: Id<"organizations">; status: unknown }, Id<"issues">>;
81+
};
6682
logs: {
6783
ingest: FunctionReference<"mutation", { deployKey: string; lines: Array<{ createdAt: number | undefined; level: "log" | "warn" | "error"; line: string }>; organizationId: Id<"organizations">; scriptName: string }, { ingested: number; }>;
6884
list: FunctionReference<"query", { afterCreatedAt?: number; organizationId: Id<"organizations">; scriptName: string }, { createdAt: number; level: "error" | "log" | "warn"; line: string; }[]>;
@@ -93,6 +109,9 @@ export interface ApiTypes {
93109
remove: FunctionReference<"mutation", { id: Id<"secrets">; organizationId: Id<"organizations"> }, void>;
94110
store: FunctionReference<"mutation", { ciphertext: string; environment?: "all" | "production" | "preview" | "dev"; iv: string; name: string; organizationId: Id<"organizations">; projectId: Id<"projects"> }, Id<"secrets">>;
95111
};
112+
telemetry: {
113+
ingest: FunctionReference<"mutation", { deployKey: string; deploymentId?: Id<"deployments">; events: Array<unknown>; organizationId: Id<"organizations"> }, { alerts: { body: string; channel: "email" | "webhook"; destination: string; id: Id<"alerts">; subject: string; }[]; incidents: number; issues: number; }>;
114+
};
96115
usage: {
97116
ingest: FunctionReference<"mutation", { deployKey: string; deploymentId?: Id<"deployments">; organizationId: Id<"organizations">; periodStart: number; quantity: number }, Id<"platformUsage">>;
98117
series: FunctionReference<"query", { organizationId: Id<"organizations">; periodStart: number }, { cpuMs: number; day: number; requests: number; }[]>;

0 commit comments

Comments
 (0)