Skip to content
Merged
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
85 changes: 85 additions & 0 deletions apps/server/src/engine/page-diff.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
/** Lines of page text for comparing two checks of a watched page. */
export function pageLines(text: string, limit = 2000): string[] {
const lines = new Set<string>();
for (const raw of text.split(/\r?\n/)) {
const line = raw.replace(/\s+/g, " ").trim().slice(0, 300);
if (line) lines.add(line);
if (lines.size >= limit) break;
}
return [...lines];
}

export type PageDiff = { added: string[]; updated: string[]; removed: string[] };

// Relative times ("posted 1 day ago" → "posted 2 days ago", "58 minutes ago" → "1 hour ago")
// tick on their own, so a change in them alone is not news.
const relativeTime =
/\b(?:\d+|an?|one)\s+(?:sec(?:ond)?|min(?:ute)?|hour|hr|day|week|month|year)s?\s+ago\b|\bvor\s+(?:\d+|einer?|einem)\s+(?:sekunde|minute|stunde|tag|woche|monat|jahr)(?:e|en|n)?\b/giu;

/** Text with relative times blanked out: equal results mean only those times changed. */
export function withoutRelativeTimes(text: string) {
return text.replace(relativeTime, "<time>");
}

// Any other number change (a price, stock count or version) is news, listed as an update.
const numberless = (line: string) =>
withoutRelativeTimes(line)
.replace(/\d+(?:[.,]\d+)*/g, "#")
.replace(/(\p{L})s\b/gu, "$1")
.toLowerCase();

/**
* Pairs each current line with one unused previous line of the same key, so one remaining
* line can never account for another that was removed.
*/
function pair(previous: string[], current: string[], key: (line: string) => string) {
const open = new Map<string, string[]>();
for (const line of previous) open.set(key(line), [...(open.get(key(line)) ?? []), line]);
const used = new Set<string>();
const paired: string[] = [];
const unpaired: string[] = [];
for (const line of current) {
const match = open.get(key(line))?.shift();
if (match === undefined) unpaired.push(line);
else {
used.add(match);
paired.push(line);
}
}
return { paired, unpaired, left: previous.filter((line) => !used.has(line)) };
}

export function diffPage(previous: string[], current: string[]): PageDiff {
const same = pair(previous, current, (line) => line);
// A line where only a relative time changed is left out.
const retimed = pair(same.left, same.unpaired, withoutRelativeTimes);
const renumbered = pair(retimed.left, retimed.unpaired, numberless);
return { added: renumbered.unpaired, updated: renumbered.paired, removed: renumbered.left };
}

/** A short "New / Updated / Removed" summary, or "" when no line changed. */
export function describePageDiff(diff: PageDiff) {
const section = (title: string, lines: string[], limit: number) => {
if (!lines.length) return [];
const shown = lines.slice(0, limit).map((line) => `• ${line.slice(0, 160)}`);
const more = lines.length > limit ? [`+${lines.length - limit} more`] : [];
return [`${title}:`, ...shown, ...more];
};
return [
...section("New", diff.added, 8),
...section("Updated", diff.updated, 4),
...section("Removed", diff.removed, 4),
].join("\n");
}

/** A one-line count for task results, such as "3 lines changed (2 new, 1 updated)". */
export function countPageDiff(diff: PageDiff) {
const total = diff.added.length + diff.updated.length + diff.removed.length;
if (!total) return "";
const parts = [
diff.added.length && `${diff.added.length} new`,
diff.updated.length && `${diff.updated.length} updated`,
diff.removed.length && `${diff.removed.length} removed`,
].filter(Boolean);
return `${total} line${total === 1 ? "" : "s"} changed (${parts.join(", ")})`;
}
45 changes: 41 additions & 4 deletions apps/server/src/engine/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,18 @@ import { SearchService } from "../search.ts";
import type { WorkspaceService } from "../workspace.ts";
import { analyzeSpending } from "./finance.ts";
import { executeModelTask } from "./model.ts";
import {
countPageDiff,
describePageDiff,
diffPage,
pageLines,
withoutRelativeTimes,
} from "./page-diff.ts";
import { LostLeaseError, type TaskContext, TaskWorker } from "./worker.ts";

const hash = (text: string) => createHash("sha256").update(text).digest("hex");
/** A change watch's last saved page: its hash, relative-time-free fingerprint and lines. */
type MonitorPage = { id: string; hash?: string; timelessHash?: string; lines: string[] };
const date = () => new Date().toISOString();
const terminal = new Set(["succeeded", "failed", "cancelled"]);
export class AgentService {
Expand Down Expand Up @@ -1076,7 +1085,23 @@ export class AgentService {
? text.toLowerCase().includes(monitor.value.toLowerCase())
: this.matchesPrice(text, Number(monitor.value));
const previouslyMatched = Boolean(task.state.matched);
const shouldNotify = matched && (monitor.condition === "change" || !previouslyMatched);
// For change watches, keep the page lines and a relative-time-free fingerprint of the whole text.
const change = monitor.condition === "change";
const lines = change ? pageLines(observation.text) : [];
const timelessHash = change ? hash(withoutRelativeTimes(text)) : undefined;
const savedPage =
matched && change
? await this.db.get<MonitorPage>(owner, "monitor-pages", monitor.id)
: undefined;
// Saved lines can run ahead of a lost task outcome; use them only for the committed baseline.
const baseline = savedPage?.hash === previousHash ? savedPage : undefined;
// Only relative times changed anywhere on the page ("3 minutes ago"): keep watching quietly.
const quiet = Boolean(baseline?.timelessHash && baseline.timelessHash === timelessHash);
const shouldNotify =
matched && !quiet && (monitor.condition === "change" || !previouslyMatched);
const diff = baseline && shouldNotify ? diffPage(baseline.lines, lines) : undefined;
// Empty when the change is past the saved lines; the alert then quotes the page instead.
const changes = diff ? describePageDiff(diff) : "";
// A notification-worthy observation is a new event, even when the page text is
// identical to an earlier one (a change back to a seen state, or a condition that
// cleared and reappeared). Commit its sequence with the outcome so publication
Expand All @@ -1100,20 +1125,30 @@ export class AgentService {
},
);
if (!savedMonitor) throw new LostLeaseError();
if (change)
await this.db.put(owner, "monitor-pages", {
id: monitor.id,
hash: currentHash,
timelessHash,
lines,
});
await ctx.event(
"observation",
previousHash ? "Checked for changes" : "Saved the first observation",
text.slice(0, 1000),
);
if (shouldNotify) {
await ctx.guard();
await ctx.event("result", "A meaningful change was found", text.slice(0, 500));
await ctx.event("result", "A meaningful change was found", changes || text.slice(0, 500));
}
const count = diff ? countPageDiff(diff) : "";
return {
status: "scheduled",
nextRunAt: nextCheckAt,
result: shouldNotify
? "Change found. A notification is ready."
? count
? `Change found: ${count}. A notification is ready.`
: "Change found. A notification is ready."
: "Watching. I'll check again on schedule.",
state: {
...task.state,
Expand All @@ -1126,7 +1161,9 @@ export class AgentService {
notice: shouldNotify
? {
title: monitor.title,
body: `Condition met at ${observation.url}: ${text.slice(0, 240)}`,
body: changes
? `Changed at ${observation.url}\n${changes}`
: `Condition met at ${observation.url}: ${text.slice(0, 240)}`,
key: `monitor:${monitor.id}:${alertSequence}:${currentHash}`,
}
: null,
Expand Down
141 changes: 140 additions & 1 deletion tests/agent-api.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import assert from "node:assert/strict";
import { mkdtemp, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { after, before, test } from "node:test";
import { after, before, beforeEach, test } from "node:test";
import { CopilotKitIntelligence } from "@copilotkit/runtime/v2";
import { createApp } from "../apps/server/src/app.ts";
import type { Config } from "../apps/server/src/config.ts";
Expand Down Expand Up @@ -55,6 +55,12 @@ before(async () => {
assert.equal(session.status, 200);
token = (await session.json()).token;
});
beforeEach(async () => {
// Keep persisted fixtures and the session, but give each independent test its
// own app middleware state, including the per-connection request budget.
await server.agent.stop();
server = await createApp(db, config);
});
after(async () => {
await server?.agent?.stop();
await db?.close();
Expand Down Expand Up @@ -449,3 +455,136 @@ test("manual milestone completion preserves current titles, ordering and concurr
404,
);
});

test("a change watch alert lists the new and updated lines of the page", async () => {
await read("/sample-page", {
text: "AI jobs in Munich\nSiemens · Werkstudent AI · 2 openings\nBCG · Intern",
});
const monitor = await read<Monitor>(
"/monitors",
{
title: "Munich AI jobs",
url: "sample://availability",
condition: "change",
intervalMinutes: 1,
},
201,
);
await server.agent.worker.tick();
await read("/sample-page", {
text: "AI jobs in Munich\nSAP · Working Student AI Engineer\nSiemens · Werkstudent AI · 3 openings\nBCG · Intern",
});
await read(`/monitors/${monitor.id}/control`, { action: "check" });
await server.agent.worker.tick();
const [alert] = (await read<AgentNotification[]>("/notifications")).filter(
(item) => item.taskId === monitor.taskId,
);
assert.equal(
alert?.body,
"Changed at sample://availability\nNew:\n• SAP · Working Student AI Engineer\nUpdated:\n• Siemens · Werkstudent AI · 3 openings",
);
const { task } = await read<{ task: AgentTask }>(`/tasks/${monitor.taskId}`);
assert.equal(
task.result,
"Change found: 2 lines changed (1 new, 1 updated). A notification is ready.",
);
await read(`/monitors/${monitor.id}/control`, { action: "stop" });
});

test("a change watch stays quiet when only relative times change", async () => {
await read("/sample-page", { text: "Jobs\nSiemens · Werkstudent AI · 3 minutes ago" });
const monitor = await read<Monitor>(
"/monitors",
{ title: "Quiet jobs", url: "sample://availability", condition: "change", intervalMinutes: 1 },
201,
);
const alerts = async () =>
(await read<AgentNotification[]>("/notifications")).filter(
(item) => item.taskId === monitor.taskId,
);
await server.agent.worker.tick();
for (const text of [
"Jobs\nSiemens · Werkstudent AI · 58 minutes ago",
"Jobs\nSiemens · Werkstudent AI · 1 hour ago",
]) {
await read("/sample-page", { text });
await read(`/monitors/${monitor.id}/control`, { action: "check" });
await server.agent.worker.tick();
}
assert.equal((await alerts()).length, 0, "ticking timestamps are not news");
await read("/sample-page", {
text: "Jobs\nSAP · Working Student AI\nSiemens · Werkstudent AI · 2 hours ago",
});
await read(`/monitors/${monitor.id}/control`, { action: "check" });
await server.agent.worker.tick();
const [alert] = await alerts();
assert.match(alert?.body ?? "", /New:\n• SAP · Working Student AI/);
await read(`/monitors/${monitor.id}/control`, { action: "stop" });
});

test("a change watch alerts when only a price changes", async () => {
await read("/sample-page", { text: "Headphones\nPrice: $399.99\nOnly 3 left" });
const monitor = await read<Monitor>(
"/monitors",
{ title: "Headphones", url: "sample://availability", condition: "change", intervalMinutes: 1 },
201,
);
await server.agent.worker.tick();
await read("/sample-page", { text: "Headphones\nPrice: $279.99\nOnly 3 left" });
await read(`/monitors/${monitor.id}/control`, { action: "check" });
await server.agent.worker.tick();
const alerts = (await read<AgentNotification[]>("/notifications")).filter(
(item) => item.taskId === monitor.taskId,
);
assert.equal(alerts.length, 1, "a price change is news");
assert.equal(alerts[0]?.body, "Changed at sample://availability\nUpdated:\n• Price: $279.99");
await read(`/monitors/${monitor.id}/control`, { action: "stop" });
});

test("a change watch alerts when a numbered listing is removed", async () => {
await read("/sample-page", { text: "Rentals\nListing 101 · 2 bed\nListing 102 · 2 bed" });
const monitor = await read<Monitor>(
"/monitors",
{ title: "Rentals", url: "sample://availability", condition: "change", intervalMinutes: 1 },
201,
);
await server.agent.worker.tick();
await read("/sample-page", { text: "Rentals\nListing 102 · 2 bed" });
await read(`/monitors/${monitor.id}/control`, { action: "check" });
await server.agent.worker.tick();
const alerts = (await read<AgentNotification[]>("/notifications")).filter(
(item) => item.taskId === monitor.taskId,
);
assert.equal(alerts.length, 1, "a removed listing is news");
assert.equal(
alerts[0]?.body,
"Changed at sample://availability\nRemoved:\n• Listing 101 · 2 bed",
);
await read(`/monitors/${monitor.id}/control`, { action: "stop" });
});

test("a change watch alerts on a change past the saved lines", async () => {
const long = "x".repeat(300);
const many = Array.from({ length: 2000 }, (_, i) => `Row ${i}`).join("\n");
for (const [before, after] of [
[`Terms\n${long} version 1`, `Terms\n${long} version 2`],
[`${many}\nLast row: open`, `${many}\nLast row: closed`],
]) {
await read("/sample-page", { text: before });
const monitor = await read<Monitor>(
"/monitors",
{ title: "Long page", url: "sample://availability", condition: "change", intervalMinutes: 1 },
201,
);
await server.agent.worker.tick();
await read("/sample-page", { text: after });
await read(`/monitors/${monitor.id}/control`, { action: "check" });
await server.agent.worker.tick();
const alerts = (await read<AgentNotification[]>("/notifications")).filter(
(item) => item.taskId === monitor.taskId,
);
assert.equal(alerts.length, 1, "a change past the saved lines is still news");
assert.match(alerts[0]?.body ?? "", /^Condition met at sample:\/\/availability: /);
await read(`/monitors/${monitor.id}/control`, { action: "stop" });
}
});
Loading
Loading