Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
26 kB · 609 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610import type { JsonObject } from "../core/json.js";import type { JazzThoughtStore } from "../jazz/store.js";import type { ThoughtEvent } from "../events/types.js";import type { AgentRun } from "../store/types.js";import { createDefaultRegistry } from "../events/registry.js";
export interface ActivityConsumerRun { id: string; agentId: string; agentVersion: number; status: AgentRun["status"]; kind: "rule" | "model"; description: string; outputCount: number;}
export interface ActivityItem { id: string; type: string; occurredAt: string; source: string; actor: string; rootEventId: string; parentEventId?: string; privacy: string; summary: string; presentation: ActivityPresentation; descendantEventCount: number; consumerRuns: ActivityConsumerRun[]; consumerRunsComplete: boolean;}
export interface ActivityPresentation { renderer: "generic" | "telegram-message" | "bluesky-post" | "email-message" | "x-post"; title: string; body: string; objectLabel: string; url?: string; urlLabel?: string; parentUrl?: string; parentLabel?: string; telegram?: { senderName: string; edited: boolean; attachmentCount: number; }; bluesky?: { uri: string; cid?: string; }; email?: { operation: "created" | "updated" | "destroyed"; senderName: string; senderEmail: string; subject: string; attachmentCount: number; }; x?: { postId: string; authorHandle?: string; createdAt?: string; replyToHandle?: string; language?: string; attachmentCount: number; possiblySensitive: boolean; };}
export interface RootActivityProjection { totalEvents: number; byType: Record<string, number>; items: ActivityItem[]; rootWindowComplete: boolean; processingPending: boolean;}
const ROOT_SOURCE_EVENT_TYPES = [...new Set([ ...createDefaultRegistry().list().map((entry) => entry.type).filter((type) => type.startsWith("stream.thought.source.")), "stream.thought.source.x.activity",])];const TERMINAL_RUN_STATUSES = new Set<AgentRun["status"]>([ "completed", "failed", "blocked", "abandoned", "skipped",]);
export async function buildRootActivity(store: JazzThoughtStore, limit = 100): Promise<RootActivityProjection> { const [events, runs] = await Promise.all([store.listEvents(), store.listRuns()]); const eventsById = new Map(events.map((event) => [event.id, event])); const descendantsByRoot = new Map<string, number>(); const runsByRoot = new Map<string, AgentRun[]>(); const byType: Record<string, number> = {}; for (const event of events) { byType[event.type] = (byType[event.type] ?? 0) + 1; if (event.id !== event.rootEventId) { descendantsByRoot.set(event.rootEventId, (descendantsByRoot.get(event.rootEventId) ?? 0) + 1); } } for (const run of runs) { const rootIds = rootIdsForRun(run, eventsById); for (const rootId of rootIds) runsByRoot.set(rootId, [...(runsByRoot.get(rootId) ?? []), run]); } const roots = events.filter((event) => ( event.id === event.rootEventId && event.type.startsWith("stream.thought.source.") )); const items = roots.slice(-limit).reverse().map((event) => ({ id: event.id, type: event.type, occurredAt: event.occurredAt, source: event.source, actor: event.actor, rootEventId: event.rootEventId, ...(event.parentEventId ? { parentEventId: event.parentEventId } : {}), privacy: event.privacy, summary: summarizeRoot(event), presentation: presentActivityEvent(event), descendantEventCount: descendantsByRoot.get(event.id) ?? 0, consumerRunsComplete: true, consumerRuns: (runsByRoot.get(event.id) ?? []).map((run) => ({ id: run.id, agentId: run.agentId, agentVersion: run.agentVersion, status: run.status, kind: runKind(run), description: describeRunResult(run, run.inputEventIds.map((id) => eventsById.get(id)).filter((input): input is ThoughtEvent => Boolean(input))), outputCount: run.outputEventIds.length, })), })); return { totalEvents: events.length, byType, items, rootWindowComplete: true, processingPending: false };}
export async function buildRecentRootObservations( store: JazzThoughtStore, limit = 100,): Promise<RootActivityProjection> { const [rootWindow, sources] = await Promise.all([ store.listRecentRootEvents({ types: ROOT_SOURCE_EVENT_TYPES, limit, maxScanned: 400 }), store.listSources(), ]); const roots = rootWindow.events; const byType: Record<string, number> = {}; for (const event of roots) byType[event.type] = (byType[event.type] ?? 0) + 1; const items = roots.slice(-limit).reverse().map((event) => ({ id: event.id, type: event.type, occurredAt: event.occurredAt, source: event.source, actor: event.actor, rootEventId: event.rootEventId, ...(event.parentEventId ? { parentEventId: event.parentEventId } : {}), privacy: event.privacy, summary: summarizeRoot(event), presentation: presentActivityEvent(event), descendantEventCount: 0, consumerRuns: [], consumerRunsComplete: false, })); return { totalEvents: sources.reduce((total, source) => total + source.lastSequence, 0), byType, items, rootWindowComplete: rootWindow.complete, processingPending: true, };}
export async function buildRecentRootActivity( store: JazzThoughtStore, limit = 100,): Promise<RootActivityProjection> { const runLimit = 250; const rootScanLimit = 400; const [rootWindow, sources, recentRuns] = await Promise.all([ store.listRecentRootEvents({ types: ROOT_SOURCE_EVENT_TYPES, limit, maxScanned: rootScanLimit }), store.listSources(), store.listRecentRuns(runLimit), ]); const roots = rootWindow.events; const rootIds = roots.map((event) => event.id); const rootEventLimit = 500; const rootEvents = await store.listEventsForRoots(rootIds, { limit: rootEventLimit }); const directTriggerIds = [...new Set(rootEvents.map((event) => event.id))]; const directRuns = await store.getRunsForTriggerEvents(directTriggerIds, { limit: runLimit }); const runs = [...new Map([...recentRuns, ...directRuns].map((run) => [run.id, run])).values()]; const triggerIds = [...new Set(runs.map((run) => run.triggerEventId))]; const triggerEvents = await store.getEvents(triggerIds); const inputsById = new Map([...rootEvents, ...roots, ...triggerEvents].map((event) => [event.id, event])); const recentRootIds = new Set(roots.map((event) => event.id)); const consumerRunsComplete = recentRuns.length < runLimit && rootEvents.length < rootEventLimit; const descendantsByRoot = new Map<string, Set<string>>(); const runsByRoot = new Map<string, AgentRun[]>(); for (const run of runs) { const inputs = run.inputEventIds.map((id) => inputsById.get(id)).filter((event): event is ThoughtEvent => Boolean(event)); const rootIds = rootIdsForRun(run, inputsById); for (const rootId of rootIds) { if (!recentRootIds.has(rootId)) continue; runsByRoot.set(rootId, [...(runsByRoot.get(rootId) ?? []), run]); const descendants = descendantsByRoot.get(rootId) ?? new Set<string>(); for (const input of inputs) if (input.id !== rootId) descendants.add(input.id); for (const outputId of run.outputEventIds) descendants.add(outputId); if (TERMINAL_RUN_STATUSES.has(run.status)) descendants.add(`terminal:${run.id}`); descendantsByRoot.set(rootId, descendants); } } const byType: Record<string, number> = {}; for (const root of roots) byType[root.type] = (byType[root.type] ?? 0) + 1; const items = roots.slice(-limit).reverse().map((event) => ({ id: event.id, type: event.type, occurredAt: event.occurredAt, source: event.source, actor: event.actor, rootEventId: event.rootEventId, ...(event.parentEventId ? { parentEventId: event.parentEventId } : {}), privacy: event.privacy, summary: summarizeRoot(event), presentation: presentActivityEvent(event), descendantEventCount: descendantsByRoot.get(event.id)?.size ?? 0, consumerRunsComplete, consumerRuns: (runsByRoot.get(event.id) ?? []).map((run) => ({ id: run.id, agentId: run.agentId, agentVersion: run.agentVersion, status: run.status, kind: runKind(run), description: describeRunResult( run, run.inputEventIds.map((id) => inputsById.get(id)).filter((input): input is ThoughtEvent => Boolean(input)), ), outputCount: run.outputEventIds.length, })), })); return { totalEvents: sources.reduce((total, source) => total + source.lastSequence, 0), byType, items, rootWindowComplete: rootWindow.complete, processingPending: false, };}
export async function buildRecentSourceActivity( store: JazzThoughtStore, runs: AgentRun[], source: string, limit = 20,): Promise<RootActivityProjection> { const sourceEvents = await store.listEvents({ source }); const roots = sourceEvents.filter((event) => event.id === event.rootEventId && event.type.startsWith("stream.thought.source.")); const inputIds = [...new Set(runs.flatMap((run) => run.inputEventIds))]; const inputEvents = await store.getEvents(inputIds); const inputsById = new Map([...sourceEvents, ...inputEvents].map((event) => [event.id, event])); const runsByRoot = new Map<string, AgentRun[]>(); for (const run of runs) { const rootIds = rootIdsForRun(run, inputsById); for (const rootId of rootIds) runsByRoot.set(rootId, [...(runsByRoot.get(rootId) ?? []), run]); } const byType: Record<string, number> = {}; for (const root of roots) byType[root.type] = (byType[root.type] ?? 0) + 1; const items = roots.slice(-limit).reverse().map((event): ActivityItem => { const consumerRuns = (runsByRoot.get(event.id) ?? []).map((run) => ({ id: run.id, agentId: run.agentId, agentVersion: run.agentVersion, status: run.status, kind: runKind(run), description: describeRunResult( run, run.inputEventIds.map((id) => inputsById.get(id)).filter((input): input is ThoughtEvent => Boolean(input)), ), outputCount: run.outputEventIds.length, })); return { id: event.id, type: event.type, occurredAt: event.occurredAt, source: event.source, actor: event.actor, rootEventId: event.rootEventId, ...(event.parentEventId ? { parentEventId: event.parentEventId } : {}), privacy: event.privacy, summary: summarizeRoot(event), presentation: presentActivityEvent(event), descendantEventCount: consumerRuns.reduce((total, run) => total + run.outputCount + 1, 0), consumerRunsComplete: true, consumerRuns, }; }); return { totalEvents: roots.length, byType, items, rootWindowComplete: true, processingPending: false };}
function rootIdsForRun(run: AgentRun, eventsById: Map<string, ThoughtEvent>): Set<string> { const rootIds = new Set<string>(); for (const inputId of new Set([run.triggerEventId, ...run.inputEventIds])) { const input = eventsById.get(inputId); if (!input) continue; rootIds.add(input.rootEventId); if (input.type !== "stream.thought.derived.event.batch" || !Array.isArray(input.payload.members)) continue; for (const reference of input.payload.members) { if (!reference || typeof reference !== "object" || Array.isArray(reference)) continue; const memberId = (reference as Record<string, unknown>).eventId; if (typeof memberId !== "string") continue; const member = eventsById.get(memberId); if (member?.type.startsWith("stream.thought.source.")) rootIds.add(member.rootEventId); } } return rootIds;}
export async function rebuildRootActivity(store: JazzThoughtStore, limit = 100): Promise<RootActivityProjection> { const projection = await buildRootActivity(store, limit); const events = await store.listEvents(); const latest = events.at(-1); await store.upsertProjection({ id: "root-activity", payload: JSON.parse(JSON.stringify(projection)) as JsonObject, lastEventId: latest?.id ?? "none", projectionVersion: 1, updatedAt: new Date().toISOString(), }); return projection;}
export function describeRunResult(run: AgentRun, inputs: ThoughtEvent[] = []): string { if (run.status === "failed") return `Failed: ${run.errorText ?? "no error detail recorded"}`; if (!run.result) return `${capitalize(run.status)}; no result recorded yet.`; const tags = stringArray(run.result.tags); const importance = typeof run.result.importance === "string" ? run.result.importance : undefined; const path = inputs.map((event) => event.payload.path).find((value): value is string => typeof value === "string"); const classification = tags.includes("specification") ? "a specification change" : tags.includes("agent-definition") ? "an agent-definition change" : tags.includes("md") ? "an ordinary Markdown file change" : undefined; const resultSummary = typeof run.result.summary === "string" ? run.result.summary : "Completed without a summary"; const firstSentence = classification ? `Classified ${path ?? "the event"} as ${classification}${importance ? ` (${importance} importance)` : ""}.` : resultSummary; const recommendation = isObject(run.result.recommendation) && typeof run.result.recommendation.proposedAction === "string" ? ` Proposed: ${run.result.recommendation.proposedAction} The proposal was not executed.` : ""; return `${firstSentence}${recommendation}`;}
function summarizeRoot(event: ThoughtEvent): string { if (typeof event.payload.path === "string") { return `${event.payload.path} ${event.type.split(".").at(-1) ?? "observed"}`; } if (event.type === "stream.thought.source.atproto.commit") { const operation = typeof event.payload.operation === "string" ? event.payload.operation : "observed"; const collection = typeof event.payload.collection === "string" ? event.payload.collection : "ATProto record"; const record = isObject(event.payload.record) ? event.payload.record : undefined; const text = record && typeof record.text === "string" ? record.text.replaceAll(/\s+/g, " ").trim() : undefined; const subject = record && isObject(record.subject) ? record.subject : undefined; const subjectUri = subject && typeof subject.uri === "string" ? subject.uri : undefined; const noun = collection === "app.bsky.feed.post" ? "Post" : collection === "app.bsky.feed.like" ? "Like" : collection; if (text) return `${capitalize(operation)} ${noun.toLowerCase()}: ${truncate(text, 180)}`; if (subjectUri) return `${capitalize(operation)} ${noun.toLowerCase()}: ${atUriToWebUrl(subjectUri)}`; return `${capitalize(operation)} ${noun}`; } if (typeof event.payload.title === "string") return event.payload.title; if (typeof event.payload.summary === "string") return event.payload.summary; return event.type;}
export function presentActivityEvent(event: ThoughtEvent): ActivityPresentation { if (event.type === "stream.thought.source.telegram.message") { const text = typeof event.payload.text === "string" ? event.payload.text : ""; const senderName = typeof event.payload.senderName === "string" && event.payload.senderName.trim() ? event.payload.senderName.trim() : "Telegram sender"; const attachments = Array.isArray(event.payload.attachments) ? event.payload.attachments : []; return { renderer: "telegram-message", title: "Telegram message", body: text, objectLabel: attachments.length === 1 ? "1 attachment" : attachments.length > 1 ? `${attachments.length} attachments` : "Message", telegram: { senderName, edited: typeof event.payload.editedAt === "string", attachmentCount: attachments.length, }, }; } if (event.type === "stream.thought.source.atproto.commit") { const payload = event.payload; const operation = typeof payload.operation === "string" ? payload.operation : "observed"; const record = isObject(payload.record) ? payload.record : undefined; const collection = typeof payload.collection === "string" ? payload.collection : undefined; const atUri = typeof payload.atUri === "string" ? payload.atUri : undefined; if (collection === "app.bsky.feed.post") { const reply = record && isObject(record.reply) && isObject(record.reply.parent) ? typeof record.reply.parent.uri === "string" ? record.reply.parent.uri : undefined : undefined; const deleted = operation === "delete"; const postUrl = atUri ? atUriToWebUrl(atUri) : ""; const parentUrl = reply ? atUriToWebUrl(reply) : ""; return { renderer: deleted || !postUrl ? "generic" : "bluesky-post", title: deleted ? "Deleted a Bluesky post" : reply ? "Replied on Bluesky" : "Posted on Bluesky", body: record && typeof record.text === "string" ? record.text : "", objectLabel: deleted ? "Post removed from Bluesky" : "Bluesky post", ...(postUrl ? { url: postUrl, urlLabel: reply ? "Open reply" : "Open post" } : {}), ...(parentUrl ? { parentUrl, parentLabel: "Open parent post" } : {}), ...(!deleted && postUrl && atUri ? { bluesky: { uri: atUri, ...(typeof payload.cid === "string" ? { cid: payload.cid } : {}), } } : {}), }; } if (collection === "app.bsky.feed.like" || collection === "app.bsky.feed.repost") { const subject = record && isObject(record.subject) && typeof record.subject.uri === "string" ? record.subject.uri : undefined; const subjectUrl = subject ? atUriToWebUrl(subject) : ""; const like = collection.endsWith("like"); const removed = operation === "delete"; return { renderer: removed || !subjectUrl ? "generic" : "bluesky-post", title: like ? removed ? "Removed a like" : "Liked a post" : removed ? "Removed a repost" : "Reposted on Bluesky", body: "", objectLabel: "Bluesky post", ...(subjectUrl ? { url: subjectUrl, urlLabel: like && !removed ? "Open liked post" : "Open post" } : {}), ...(!removed && subjectUrl && subject ? { bluesky: { uri: subject, ...(record && isObject(record.subject) && typeof record.subject.cid === "string" ? { cid: record.subject.cid } : {}), } } : {}), }; } if (collection === "app.bsky.graph.follow") { const subject = record && typeof record.subject === "string" ? record.subject : undefined; const removed = operation === "delete"; return { renderer: "generic", title: removed ? "Stopped following an account" : "Followed an account", body: "", objectLabel: "Bluesky account", ...(subject ? { url: `https://bsky.app/profile/${encodeURIComponent(subject)}`, urlLabel: "Open profile" } : {}), }; } } if (event.type === "stream.thought.source.x.activity") { const eventType = typeof event.payload.eventType === "string" ? event.payload.eventType : ""; const post = isObject(event.payload.post) ? event.payload.post : undefined; const like = isObject(event.payload.like) ? event.payload.like : undefined; const postId = typeof post?.postId === "string" ? post.postId : typeof like?.postId === "string" ? like.postId : undefined; const postUrl = postId ? `https://x.com/i/web/status/${encodeURIComponent(postId)}` : ""; if (eventType === "post.create") { const reply = typeof post?.inReplyToUserId === "string" || (Array.isArray(post?.referencedPosts) && post.referencedPosts.some((reference) => ( isObject(reference) && reference.type === "replied_to" ))); const authorHandle = xHandleFromSubscriptionTag(event); const replyToHandle = xReplyToHandle(post); const replyPostId = xReferencedPostId(post, "replied_to"); const quotedPostId = xReferencedPostId(post, "quoted"); const relatedPostId = replyPostId ?? quotedPostId; const parentUrl = relatedPostId ? `https://x.com/i/web/status/${encodeURIComponent(relatedPostId)}` : ""; const text = typeof post?.text === "string" ? post.text : ""; const attachments = isObject(post?.attachments) && Array.isArray(post.attachments.mediaKeys) ? post.attachments.mediaKeys.length : 0; return { renderer: "x-post", title: reply ? "Replied on X" : "Posted on X", body: xDisplayText(text, replyToHandle), objectLabel: "X post", ...(postUrl ? { url: postUrl, urlLabel: reply ? "Open reply" : "Open post" } : {}), ...(parentUrl ? { parentUrl, parentLabel: replyPostId ? "Open parent post" : "Open quoted post", } : {}), ...(postId ? { x: { postId, ...(authorHandle ? { authorHandle } : {}), ...(typeof post?.createdAt === "string" ? { createdAt: post.createdAt } : {}), ...(replyToHandle ? { replyToHandle } : {}), ...(typeof post?.language === "string" ? { language: post.language } : {}), attachmentCount: attachments, possiblySensitive: post?.possiblySensitive === true, } } : {}), }; } if (eventType === "post.delete") { return { renderer: "generic", title: "Deleted an X post", body: "", objectLabel: "X post", }; } if (eventType === "like.create") { return { renderer: "generic", title: "Liked a post on X", body: "", objectLabel: "X post", ...(postUrl ? { url: postUrl, urlLabel: "Open liked post" } : {}), }; } } if (event.type === "stream.thought.source.email.observed") { const operation = event.payload.operation === "created" || event.payload.operation === "destroyed" ? event.payload.operation : "updated"; const subject = typeof event.payload.subject === "string" ? event.payload.subject : ""; const preview = typeof event.payload.preview === "string" ? event.payload.preview : ""; const sender = Array.isArray(event.payload.from) && isObject(event.payload.from[0]) ? event.payload.from[0] : undefined; const senderName = sender && typeof sender.name === "string" ? sender.name : ""; const senderEmail = sender && typeof sender.email === "string" ? sender.email : ""; const attachmentCount = Array.isArray(event.payload.attachments) ? event.payload.attachments.length : 0; return { renderer: "email-message", title: operation === "created" ? "New email" : operation === "destroyed" ? "Email removed" : "Email updated", body: operation === "destroyed" ? "" : preview, objectLabel: subject || (operation === "destroyed" ? "Message removed" : "Email without a subject"), email: { operation, senderName, senderEmail, subject, attachmentCount }, }; } const summary = summarizeRoot(event); return { renderer: "generic", title: capitalize(event.type.split(".").at(-1) ?? "observation"), body: summary === event.type ? "" : summary, objectLabel: "Observation", };}
function runKind(run: AgentRun): ActivityConsumerRun["kind"] { return run.provider === "deterministic" && run.model === "deterministic" ? "rule" : "model";}
function stringArray(value: unknown): string[] { return Array.isArray(value) ? value.filter((item): item is string => typeof item === "string") : [];}
function isObject(value: unknown): value is JsonObject { return Boolean(value && typeof value === "object" && !Array.isArray(value));}
function capitalize(value: string): string { return value.charAt(0).toUpperCase() + value.slice(1);}
function truncate(value: string, limit: number): string { return value.length <= limit ? value : `${value.slice(0, limit - 1).trimEnd()}…`;}
function atUriToWebUrl(value: string): string { const match = value.match(/^at:\/\/([^/]+)\/app\.bsky\.feed\.post\/([^/]+)$/); return match ? `https://bsky.app/profile/${blueskyActorPathSegment(match[1]!)}/post/${encodeURIComponent(match[2]!)}` : "";}
function blueskyActorPathSegment(value: string): string { return encodeURIComponent(value).replaceAll("%3A", ":");}
function xHandleFromSubscriptionTag(event: ThoughtEvent): string | undefined { const tag = typeof event.payload.subscriptionTag === "string" ? event.payload.subscriptionTag : ""; const eventType = typeof event.payload.eventType === "string" ? event.payload.eventType : ""; const prefix = `thoughtstream:${event.source}:`; const suffix = `:${eventType.replaceAll(".", "-")}`; if (!tag.startsWith(prefix) || !tag.endsWith(suffix)) return undefined; const handle = tag.slice(prefix.length, -suffix.length); return /^[A-Za-z0-9_]{1,50}$/.test(handle) ? handle : undefined;}
function xReplyToHandle(post: JsonObject | undefined): string | undefined { const replyUserId = typeof post?.inReplyToUserId === "string" ? post.inReplyToUserId : undefined; const entities = isObject(post?.entities) ? post.entities : undefined; const mentions = entities && Array.isArray(entities.mentions) ? entities.mentions : []; for (const mention of mentions) { if (!isObject(mention) || typeof mention.username !== "string") continue; if (replyUserId && mention.userId !== replyUserId) continue; if (/^[A-Za-z0-9_]{1,50}$/.test(mention.username)) return mention.username; } return undefined;}
function xReferencedPostId(post: JsonObject | undefined, type: "replied_to" | "quoted"): string | undefined { const references = Array.isArray(post?.referencedPosts) ? post.referencedPosts : []; for (const reference of references) { if (isObject(reference) && reference.type === type && typeof reference.postId === "string") { return reference.postId; } } return undefined;}
function xDisplayText(text: string, replyToHandle: string | undefined): string { if (!replyToHandle) return text; const prefix = `@${replyToHandle}`; return text.slice(0, prefix.length).toLowerCase() === prefix.toLowerCase() ? text.slice(prefix.length).trimStart() : text;}