diff --git a/apps/checker/cmd/main.go b/apps/checker/cmd/main.go index 2d00178c..4f587805 100644 --- a/apps/checker/cmd/main.go +++ b/apps/checker/cmd/main.go @@ -107,6 +107,7 @@ func main() { StatusCode: res.StatusCode, Region: flyRegion, Message: res.Message, + CronTimestamp: req.CronTimestamp, }) } else if req.Status == "error" && statusCode.IsSuccessful() { // Q: Why here we check the data before updating the status in this scenario? @@ -115,6 +116,8 @@ func main() { Status: "active", Region: flyRegion, StatusCode: res.StatusCode, + CronTimestamp: req.CronTimestamp, + }) } diff --git a/apps/checker/ping.go b/apps/checker/ping.go index 538544e8..5dbc3b72 100644 --- a/apps/checker/ping.go +++ b/apps/checker/ping.go @@ -100,9 +100,9 @@ func Ping(ctx context.Context, client *http.Client, inputData request.CheckerReq latency := time.Since(start).Milliseconds() if err != nil { - timingAsString, err := json.Marshal(timing) - if err != nil { - logger.Error().Err(err).Msg("error while parsing timing data") + timingAsString, err2 := json.Marshal(timing) + if err2 != nil { + logger.Error().Err(err2).Msg("error while parsing timing data") } var urlErr *url.Error if errors.As(err, &urlErr) && urlErr.Timeout() { diff --git a/apps/checker/update.go b/apps/checker/update.go index 47b2c19c..a8146300 100644 --- a/apps/checker/update.go +++ b/apps/checker/update.go @@ -12,11 +12,12 @@ import ( ) type UpdateData struct { - MonitorId string `json:"monitorId"` - Status string `json:"status"` - Message string `json:"message,omitempty"` - StatusCode int `json:"statusCode,omitempty"` - Region string `json:"region"` + MonitorId string `json:"monitorId"` + Status string `json:"status"` + Message string `json:"message,omitempty"` + StatusCode int `json:"statusCode,omitempty"` + Region string `json:"region"` + CronTimestamp int64 `json:"cronTimestamp"` } func UpdateStatus(ctx context.Context, updateData UpdateData) { diff --git a/apps/server/src/checker/alerting.test.ts b/apps/server/src/checker/alerting.test.ts index ce79429e..f2f5e07e 100644 --- a/apps/server/src/checker/alerting.test.ts +++ b/apps/server/src/checker/alerting.test.ts @@ -11,6 +11,10 @@ test.todo("should send email notification", async () => { }, }; }); - await triggerAlerting({ monitorId: "1", region: "ams", statusCode: 400 }); + await triggerAlerting({ + monitorId: "1", + region: "ams", + statusCode: 400, + }); expect(fn).toHaveBeenCalled(); }); diff --git a/apps/server/src/checker/index.ts b/apps/server/src/checker/index.ts index 0de6bd20..e529cc87 100644 --- a/apps/server/src/checker/index.ts +++ b/apps/server/src/checker/index.ts @@ -1,13 +1,18 @@ import { Hono } from "hono"; import { z } from "zod"; +import { eq, schema } from "@openstatus/db"; +import { db } from "@openstatus/db/src/db"; import { flyRegions } from "@openstatus/db/src/schema/monitors/constants"; +import { selectMonitorSchema } from "@openstatus/db/src/schema/monitors/validation"; +import { Redis } from "@openstatus/upstash"; import { env } from "../env"; import { checkerAudit } from "../utils/audit-log"; import { triggerAlerting, upsertMonitorStatus } from "./alerting"; export const checkerRoute = new Hono(); +const redis = Redis.fromEnv(); checkerRoute.post("/updateStatus", async (c) => { const auth = c.req.header("Authorization"); @@ -17,20 +22,22 @@ checkerRoute.post("/updateStatus", async (c) => { } const json = await c.req.json(); - const schema = z.object({ + const payloadSchema = z.object({ monitorId: z.string(), status: z.enum(["active", "error"]), // that's the new status message: z.string().optional(), statusCode: z.number().optional(), region: z.enum(flyRegions), + cronTimestamp: z.number().optional(), }); - const result = schema.safeParse(json); + const result = payloadSchema.safeParse(json); if (!result.success) { // console.error(result.error); return c.text("Unprocessable Entity", 422); } - const { monitorId, status, message, region, statusCode } = result.data; + const { monitorId, status, message, region, statusCode, cronTimestamp } = + result.data; console.log(`📝 update monitor status ${JSON.stringify(result.data)}`); @@ -49,10 +56,14 @@ checkerRoute.post("/updateStatus", async (c) => { id: `monitor:${monitorId}`, action: "monitor.recovered", targets: [{ id: monitorId, type: "monitor" }], - metadata: { region: region, statusCode: statusCode }, + metadata: { + region: region, + statusCode: statusCode, + cronTimestamp: cronTimestamp, + }, }); break; - case "error": + case "error": { await upsertMonitorStatus({ monitorId: monitorId, status: "error", @@ -67,8 +78,38 @@ checkerRoute.post("/updateStatus", async (c) => { region: region, statusCode: statusCode, message, + cronTimestamp, }, }); + + const currentMonitor = await db + .select() + .from(schema.monitor) + .where(eq(schema.monitor.id, Number(monitorId))) + .get(); + + const monitor = selectMonitorSchema.parse(currentMonitor); + + if (!cronTimestamp) { + console.log("cronTimestamp is undefined"); + } + + const redisKey = `${monitorId}-${cronTimestamp}`; + // We add the new region to the set + await redis.sadd(redisKey, region); + // let's add an expire to the set + await redis.expire(redisKey, 60 * 60 * 24); + // We get the number of regions affected + const nbAffectedRegion = await redis.scard(redisKey); + + // If the number of affected regions is greater than half of the total region, we trigger the alerting + if (nbAffectedRegion < monitor.regions.length / 2) { + console.log( + `Not enough affected regions (${nbAffectedRegion}/${flyRegions.length})`, + ); + break; + } + // Add a await triggerAlerting({ monitorId: monitorId, region: env.FLY_REGION, @@ -76,6 +117,7 @@ checkerRoute.post("/updateStatus", async (c) => { message, }); break; + } } return c.text("Ok", 200); }); diff --git a/packages/tinybird/src/audit-log/action-schema.ts b/packages/tinybird/src/audit-log/action-schema.ts index cf1c00c5..6d710f35 100644 --- a/packages/tinybird/src/audit-log/action-schema.ts +++ b/packages/tinybird/src/audit-log/action-schema.ts @@ -6,7 +6,11 @@ import { z } from "zod"; */ export const monitorRecoveredSchema = z.object({ action: z.literal("monitor.recovered"), - metadata: z.object({ region: z.string(), statusCode: z.number() }), + metadata: z.object({ + region: z.string(), + statusCode: z.number(), + cronTimestamp: z.number().optional(), + }), }); /** @@ -19,6 +23,7 @@ export const monitorFailedSchema = z.object({ region: z.string(), statusCode: z.number().optional(), message: z.string().optional(), + cronTimestamp: z.number().optional(), }), });