You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
472 lines
18 KiB
472 lines
18 KiB
#!/usr/bin/env node
|
|
// scrape-channel.mjs — Fetch Teams channel messages without ChannelMessage.Read.All.
|
|
//
|
|
// Captures a bearer token from the user's signed-in Teams web session (Playwright with
|
|
// persistent profile + Edge channel), then replays paginated GETs against Teams'
|
|
// internal CSA (Conversation Service Adapter) endpoint to retrieve channel messages
|
|
// in a Graph-compatible shape.
|
|
//
|
|
// USAGE
|
|
// node scrape-channel.mjs --team-id <guid> --channel-id <19:...@thread.tacv2>
|
|
// [--since ISO] [--until ISO] [--limit N]
|
|
// [--headed] [--refresh-token] [--include-replies]
|
|
//
|
|
// OUTPUT
|
|
// JSON to stdout: { value: [ message, ... ], capturedAt, source }
|
|
// Each message: { id, createdDateTime, lastModifiedDateTime, from, body, attachments,
|
|
// replyToId, replies?: [ ... ] }
|
|
//
|
|
// EXIT CODES
|
|
// 0 ok, 2 bad args, 3 token-capture failed, 4 fetch failed after retry.
|
|
|
|
import { chromium } from "playwright";
|
|
import { mkdirSync, existsSync, readFileSync, writeFileSync } from "node:fs";
|
|
import { dirname, join } from "node:path";
|
|
import { fileURLToPath } from "node:url";
|
|
|
|
const __dirname = dirname(fileURLToPath(import.meta.url));
|
|
const PROFILE_DIR = join(__dirname, ".pw-profile");
|
|
const TOKEN_CACHE = join(__dirname, ".token-cache.json");
|
|
const CAPTURE_TIMEOUT_MS = Number.parseInt(process.env.SCRAPE_CAPTURE_TIMEOUT_MS || "300000", 10);
|
|
const TOKEN_SAFETY_MARGIN_S = 120;
|
|
|
|
// --- arg parsing -----------------------------------------------------------
|
|
function parseArgs(argv) {
|
|
const out = { headed: false, refreshToken: false, includeReplies: true, limit: 0 };
|
|
for (let i = 2; i < argv.length; i++) {
|
|
const a = argv[i];
|
|
const next = () => argv[++i];
|
|
switch (a) {
|
|
case "--team-id": out.teamId = next(); break;
|
|
case "--channel-id": out.channelId = next(); break;
|
|
case "--since": out.since = next(); break;
|
|
case "--until": out.until = next(); break;
|
|
case "--limit": out.limit = Number.parseInt(next(), 10) || 0; break;
|
|
case "--headed": out.headed = true; break;
|
|
case "--refresh-token": out.refreshToken = true; break;
|
|
case "--include-replies": out.includeReplies = true; break;
|
|
case "--no-replies": out.includeReplies = false; break;
|
|
case "-h": case "--help": out.help = true; break;
|
|
default: throw new Error(`Unknown arg: ${a}`);
|
|
}
|
|
}
|
|
return out;
|
|
}
|
|
|
|
const args = parseArgs(process.argv);
|
|
if (args.help) {
|
|
process.stderr.write(`See header of ${import.meta.url}\n`);
|
|
process.exit(0);
|
|
}
|
|
if (!args.teamId || !args.channelId) {
|
|
process.stderr.write("ERR: --team-id and --channel-id are required\n");
|
|
process.exit(2);
|
|
}
|
|
|
|
const sinceMs = args.since ? Date.parse(args.since) : 0;
|
|
const untilMs = args.until ? Date.parse(args.until) : Number.POSITIVE_INFINITY;
|
|
if (Number.isNaN(sinceMs) || Number.isNaN(untilMs)) {
|
|
process.stderr.write("ERR: --since/--until must be ISO timestamps\n");
|
|
process.exit(2);
|
|
}
|
|
|
|
// --- token cache -----------------------------------------------------------
|
|
function loadCache() {
|
|
if (!existsSync(TOKEN_CACHE)) return null;
|
|
try {
|
|
const c = JSON.parse(readFileSync(TOKEN_CACHE, "utf8"));
|
|
if (!c.token || !c.expiresAt) return null;
|
|
if (Date.now() / 1000 > c.expiresAt - TOKEN_SAFETY_MARGIN_S) return null;
|
|
return c;
|
|
} catch { return null; }
|
|
}
|
|
function saveCache(c) { writeFileSync(TOKEN_CACHE, JSON.stringify(c, null, 2)); }
|
|
|
|
function decodeJwtExp(jwt) {
|
|
try {
|
|
const payload = JSON.parse(Buffer.from(jwt.split(".")[1], "base64url").toString("utf8"));
|
|
return payload.exp || 0;
|
|
} catch { return 0; }
|
|
}
|
|
|
|
// --- token capture via Playwright -----------------------------------------
|
|
async function captureToken({ teamId, channelId, headed }) {
|
|
mkdirSync(PROFILE_DIR, { recursive: true });
|
|
const ctx = await chromium.launchPersistentContext(PROFILE_DIR, {
|
|
channel: "msedge",
|
|
headless: !headed,
|
|
viewport: { width: 1280, height: 900 },
|
|
});
|
|
|
|
let resolveCap, rejectCap;
|
|
const capPromise = new Promise((res, rej) => { resolveCap = res; rejectCap = rej; });
|
|
|
|
// Accept any URL that fetches messages for our specific channel/thread id.
|
|
// Teams v2 routes vary by region and feature flags; we don't try to predict the
|
|
// host, we just look for the channel id + a /messages segment with a bearer token.
|
|
const channelIdLower = channelId.toLowerCase();
|
|
const channelIdEnc = encodeURIComponent(channelId).toLowerCase();
|
|
const matchers = [
|
|
(u) => {
|
|
const ul = u.toLowerCase();
|
|
const hasChannel = ul.includes(channelIdLower) || ul.includes(channelIdEnc);
|
|
const hasMessages = /\/messages(\/|\?|$)/i.test(ul);
|
|
return hasChannel && hasMessages;
|
|
},
|
|
// Fallback: NG chat service for the channel, even if id is opaque in URL
|
|
(u) => /\.ng\.msg\.teams\.(microsoft|live)\.com\/v\d+\/users\/ME\/conversations\/[^/]+\/messages/i.test(u),
|
|
(u) => /api\.flightproxy\.teams\.microsoft\.com\/api\/v\d+\/ep\/[^/]+\/v\d+\/users\/ME\/conversations\/[^/]+\/messages/i.test(u),
|
|
];
|
|
|
|
const seenUrls = new Set();
|
|
const onRequest = (req) => {
|
|
const url = req.url();
|
|
// Light-touch debug log so we can see what's flying past if capture fails.
|
|
if (process.env.SCRAPE_DEBUG && /messages|conversations\//i.test(url) && /teams|skype|flightproxy|trouter/i.test(url)) {
|
|
if (!seenUrls.has(url)) {
|
|
seenUrls.add(url);
|
|
process.stderr.write(`[req] ${req.method()} ${url}\n`);
|
|
}
|
|
}
|
|
const matchIdx = matchers.findIndex((m) => m(url));
|
|
if (matchIdx < 0) return;
|
|
const headers = req.headers();
|
|
const auth = headers["authorization"];
|
|
if (!auth || !auth.startsWith("Bearer ")) return;
|
|
const token = auth.slice(7);
|
|
const exp = decodeJwtExp(token);
|
|
const u = new URL(url);
|
|
process.stderr.write(`Captured token (matcher #${matchIdx + 1}) from ${u.origin}${u.pathname}\n`);
|
|
resolveCap({
|
|
token,
|
|
expiresAt: exp || (Math.floor(Date.now() / 1000) + 3000),
|
|
endpoint: `${u.origin}${u.pathname.replace(/\/messages.*/, "/messages")}`,
|
|
headers: {
|
|
accept: headers["accept"] || "application/json",
|
|
"accept-language": headers["accept-language"] || "en-US,en;q=0.9",
|
|
"user-agent": headers["user-agent"] || "",
|
|
// Pass through any x-ms-* / behavior headers the service may require.
|
|
...Object.fromEntries(
|
|
Object.entries(headers).filter(([k]) =>
|
|
k.startsWith("x-ms-") || k === "behavioroverride" || k === "x-skypetoken",
|
|
),
|
|
),
|
|
},
|
|
capturedAt: new Date().toISOString(),
|
|
});
|
|
};
|
|
ctx.on("request", onRequest);
|
|
|
|
const groupId = teamId;
|
|
const encodedChannel = encodeURIComponent(channelId);
|
|
// New Teams ("v2") deep-link variants. Try in order; first non-erroring wins.
|
|
// If none land directly on the channel, the user can click into it manually
|
|
// inside the Playwright window; the interceptor still fires.
|
|
const deepLinks = [
|
|
`https://teams.microsoft.com/v2/#/channel/${encodedChannel}/conversations?groupId=${groupId}`,
|
|
`https://teams.microsoft.com/v2/#/conversations/channel?threadId=${encodedChannel}&ctx=channel&groupId=${groupId}`,
|
|
`https://teams.microsoft.com/l/channel/${encodedChannel}/General?groupId=${groupId}`,
|
|
`https://teams.microsoft.com/v2/`,
|
|
`https://teams.microsoft.com/`,
|
|
];
|
|
const page = await ctx.newPage();
|
|
|
|
let timer;
|
|
try {
|
|
process.stderr.write(
|
|
`Opening Teams. If sign-in is required, complete it in the browser window. `
|
|
+ `Up to ${Math.round(CAPTURE_TIMEOUT_MS / 1000)}s to capture a channel request.\n`,
|
|
);
|
|
for (const link of deepLinks) {
|
|
try {
|
|
await page.goto(link, { waitUntil: "domcontentloaded", timeout: 30_000 });
|
|
break;
|
|
} catch { /* try next */ }
|
|
}
|
|
timer = setTimeout(
|
|
() => rejectCap(new Error(
|
|
`token capture timed out after ${Math.round(CAPTURE_TIMEOUT_MS / 1000)}s. `
|
|
+ `Re-run with SCRAPE_DEBUG=1 and --headed; manually click into the channel if needed.`,
|
|
)),
|
|
CAPTURE_TIMEOUT_MS,
|
|
);
|
|
const result = await capPromise;
|
|
return result;
|
|
} finally {
|
|
clearTimeout(timer);
|
|
await ctx.close().catch(() => {});
|
|
}
|
|
}
|
|
|
|
// --- message fetch ---------------------------------------------------------
|
|
async function fetchPage(cache, url) {
|
|
const res = await fetch(url, {
|
|
headers: {
|
|
...cache.headers,
|
|
authorization: `Bearer ${cache.token}`,
|
|
},
|
|
});
|
|
return res;
|
|
}
|
|
|
|
function isInWindow(msg) {
|
|
const t = Date.parse(msg.createdDateTime || msg.composetime || msg.originalArrivalTime || 0);
|
|
if (Number.isNaN(t)) return true;
|
|
return t >= sinceMs && t <= untilMs;
|
|
}
|
|
|
|
function parseFromField(rawFrom, imdisplayname) {
|
|
if (imdisplayname && typeof imdisplayname === "string" && !imdisplayname.startsWith("http")) {
|
|
return { displayName: imdisplayname, userId: null };
|
|
}
|
|
if (rawFrom && typeof rawFrom === "object" && rawFrom.user) {
|
|
return { displayName: rawFrom.user.displayName || null, userId: rawFrom.user.id || null };
|
|
}
|
|
if (typeof rawFrom === "string") {
|
|
// Skype/Teams contact URL: .../contacts/<id> where id is either
|
|
// "8:orgid:<guid>" (real user) or "19:<thread>@thread.tacv2" (channel itself).
|
|
const m = /\/contacts\/(8:orgid:[^/?#]+|19:[^/?#]+)/i.exec(rawFrom);
|
|
if (m) {
|
|
const id = decodeURIComponent(m[1]);
|
|
if (id.startsWith("19:")) return { displayName: "Channel", userId: id };
|
|
if (id.startsWith("8:orgid:")) {
|
|
const guid = id.slice("8:orgid:".length);
|
|
return { displayName: null, userId: guid };
|
|
}
|
|
return { displayName: null, userId: id };
|
|
}
|
|
}
|
|
return { displayName: null, userId: null };
|
|
}
|
|
|
|
// Some CSA responses return Teams chat-service shape; normalize to Graph-ish.
|
|
function normalize(msg) {
|
|
if (msg.body && typeof msg.body === "object") return msg; // already Graph-shaped
|
|
const created = msg.composetime || msg.originalarrivaltime || msg.originalArrivalTime || msg.createdDateTime;
|
|
const contentType = (msg.messagetype || "").toLowerCase().includes("html") ? "html" : "text";
|
|
const { displayName, userId } = parseFromField(msg.from, msg.imdisplayname);
|
|
// Channel replies carry their root id in the conversationLink as
|
|
// `...;messageid=<rootId>`. Root posts have a conversationLink with no
|
|
// messageid suffix (or messageid equal to their own id).
|
|
const link = msg.conversationLink || msg.conversationid || "";
|
|
const m = /messageid=(\d+)/i.exec(link);
|
|
const rootId = m ? m[1] : null;
|
|
const ownId = String(msg.id || "");
|
|
const parent = rootId && rootId !== ownId ? rootId : null;
|
|
let attachments = [];
|
|
const props = msg.properties || {};
|
|
if (props.files) {
|
|
try { attachments = typeof props.files === "string" ? JSON.parse(props.files) : props.files; }
|
|
catch { attachments = []; }
|
|
}
|
|
return {
|
|
id: msg.id || msg.clientmessageid,
|
|
createdDateTime: created,
|
|
lastModifiedDateTime: msg.version || created,
|
|
from: { user: { displayName: displayName, id: userId } },
|
|
body: { contentType, content: msg.content || "" },
|
|
attachments,
|
|
replyToId: parent,
|
|
messageType: msg.messagetype || null,
|
|
subject: props.subject || null,
|
|
};
|
|
}
|
|
|
|
// Resolve user GUIDs to display names via Teams' profile lookup endpoint.
|
|
// Best-effort: failures leave the GUID in place and don't fail the run.
|
|
async function resolveDisplayNames(cache, flat, { teamId, channelId }) {
|
|
const unresolved = new Map(); // guid -> [message refs]
|
|
for (const m of flat) {
|
|
const u = m.from?.user;
|
|
if (u && !u.displayName && u.id && /^[0-9a-f]{8}-[0-9a-f]{4}-/i.test(u.id)) {
|
|
if (!unresolved.has(u.id)) unresolved.set(u.id, []);
|
|
unresolved.get(u.id).push(m);
|
|
}
|
|
}
|
|
if (unresolved.size === 0) return;
|
|
|
|
// teams.cloud.microsoft profile endpoint accepts a small batch of userIds via
|
|
// a usersInfo query param. We chunk into groups of 20 to stay polite.
|
|
const ids = [...unresolved.keys()];
|
|
const chunkSize = 20;
|
|
// Derive a thread context required by the endpoint. Use the channel id itself.
|
|
const threadId = channelId;
|
|
const base = `https://teams.cloud.microsoft/api/mt/part/msft/beta/users/me/threads/`
|
|
+ `${encodeURIComponent(threadId)}/properties/pictureV2`;
|
|
for (let i = 0; i < ids.length; i += chunkSize) {
|
|
const chunk = ids.slice(i, i + chunkSize);
|
|
const usersInfo = chunk.map((g) => ({
|
|
userId: `8:orgid:${g}`,
|
|
displayName: "",
|
|
avatarETag: "",
|
|
}));
|
|
const url = `${base}?usersInfo=${encodeURIComponent(JSON.stringify(usersInfo))}`
|
|
+ `&size=HR64x64`;
|
|
try {
|
|
const res = await fetchPage(cache, url);
|
|
if (!res.ok) continue;
|
|
const data = await res.json();
|
|
// Response shape varies; look for any displayName fields keyed by user id.
|
|
const flatten = JSON.stringify(data);
|
|
for (const g of chunk) {
|
|
// Heuristic: find `"displayName":"X"` near the user GUID in the response.
|
|
const re = new RegExp(
|
|
`"userId"\\s*:\\s*"8:orgid:${g}"[^}]*?"displayName"\\s*:\\s*"([^"]+)"`
|
|
+ `|"displayName"\\s*:\\s*"([^"]+)"[^}]*?"userId"\\s*:\\s*"8:orgid:${g}"`,
|
|
"i",
|
|
);
|
|
const m = re.exec(flatten);
|
|
const name = m ? (m[1] || m[2]) : null;
|
|
if (name) {
|
|
for (const msg of unresolved.get(g)) {
|
|
msg.from.user.displayName = name;
|
|
}
|
|
}
|
|
}
|
|
} catch { /* best-effort */ }
|
|
}
|
|
// For any still-unresolved, fall back to a short form so output isn't empty.
|
|
for (const [g, msgs] of unresolved) {
|
|
for (const m of msgs) {
|
|
if (!m.from.user.displayName) m.from.user.displayName = `User ${g.slice(0, 8)}`;
|
|
}
|
|
}
|
|
}
|
|
|
|
|
|
// /messages?$expand=replies shape. Replies are sorted oldest-first within each
|
|
// root; roots remain in the order the chat service returned them (newest first).
|
|
function groupReplies(flat) {
|
|
const byId = new Map();
|
|
for (const m of flat) byId.set(m.id, m);
|
|
const roots = [];
|
|
const orphans = [];
|
|
for (const m of flat) {
|
|
if (!m.replyToId) {
|
|
m.replies = [];
|
|
roots.push(m);
|
|
}
|
|
}
|
|
for (const m of flat) {
|
|
if (!m.replyToId) continue;
|
|
const parent = byId.get(m.replyToId);
|
|
if (parent) {
|
|
parent.replies = parent.replies || [];
|
|
parent.replies.push(m);
|
|
} else {
|
|
// Parent fell outside our window. Promote the reply to a root so it isn't
|
|
// silently dropped; flag it so the renderer can mark it as a fragment.
|
|
m.replies = [];
|
|
m.isOrphanReply = true;
|
|
orphans.push(m);
|
|
}
|
|
}
|
|
for (const r of roots) {
|
|
r.replies.sort((a, b) => Date.parse(a.createdDateTime) - Date.parse(b.createdDateTime));
|
|
}
|
|
// Filter out non-message system events (member adds, topic changes, etc.) at
|
|
// the root level — keep them only if they're replies the user might want context for.
|
|
const isContent = (m) => {
|
|
const t = (m.messageType || "").toLowerCase();
|
|
return !t || t.startsWith("text") || t.startsWith("richtext");
|
|
};
|
|
return [...roots.filter(isContent), ...orphans.filter(isContent)];
|
|
}
|
|
|
|
async function fetchMessages(cache, { teamId, channelId, limit }) {
|
|
let url = `${cache.endpoint}?pageSize=50`;
|
|
if (cache.endpoint.includes("/api/csa/")) {
|
|
// CSA endpoint already references the channel via path; nothing to add.
|
|
} else {
|
|
// ng.msg endpoint — startTime narrows server-side. Teams chatsvc expects
|
|
// epoch milliseconds, not ISO strings.
|
|
if (args.since) url += `&startTime=${sinceMs}`;
|
|
}
|
|
|
|
const collected = [];
|
|
const collectedIds = new Set();
|
|
// Roots that pass the window filter — used to know when we can stop. Replies
|
|
// for those roots can appear later in pagination even if their own timestamp
|
|
// sits outside the requested window, so we keep collecting until we run out
|
|
// of pages or until ALL their replies have been seen.
|
|
const rootIdsInWindow = new Set();
|
|
let pageCount = 0;
|
|
let stopOnNextPage = false;
|
|
|
|
while (url) {
|
|
pageCount++;
|
|
let res = await fetchPage(cache, url);
|
|
if (res.status === 401) {
|
|
process.stderr.write("Token expired mid-fetch; refreshing...\n");
|
|
const fresh = await captureToken({ teamId, channelId, headed: args.headed });
|
|
saveCache(fresh);
|
|
cache = fresh;
|
|
res = await fetchPage(cache, url);
|
|
}
|
|
if (!res.ok) {
|
|
const body = await res.text().catch(() => "");
|
|
throw new Error(`HTTP ${res.status} for ${url}\n${body.slice(0, 500)}`);
|
|
}
|
|
const data = await res.json();
|
|
const items = (data.value || data.messages || []).map(normalize);
|
|
let sawInWindowOnThisPage = false;
|
|
for (const m of items) {
|
|
if (collectedIds.has(m.id)) continue;
|
|
const inWindow = isInWindow(m);
|
|
const isReplyToTrackedRoot = m.replyToId && rootIdsInWindow.has(m.replyToId);
|
|
if (!inWindow && !isReplyToTrackedRoot) continue;
|
|
collected.push(m);
|
|
collectedIds.add(m.id);
|
|
if (!m.replyToId) rootIdsInWindow.add(m.id);
|
|
if (inWindow) sawInWindowOnThisPage = true;
|
|
if (limit && collected.filter((x) => !x.replyToId).length >= limit) {
|
|
stopOnNextPage = true;
|
|
}
|
|
}
|
|
// Heuristic stop: once a full page brought zero in-window items AND we have
|
|
// enough roots, assume earlier pages won't help.
|
|
if (stopOnNextPage && !sawInWindowOnThisPage) break;
|
|
url = data["@odata.nextLink"] || data._metadata?.syncState || data.backwardLink || null;
|
|
if (pageCount > 200) break; // hard safety stop
|
|
}
|
|
return args.includeReplies === false ? collected.filter((m) => !m.replyToId) : collected;
|
|
}
|
|
|
|
// --- main ------------------------------------------------------------------
|
|
(async () => {
|
|
let cache = args.refreshToken ? null : loadCache();
|
|
if (!cache) {
|
|
process.stderr.write(`Capturing Teams web token (profile: ${PROFILE_DIR})...\n`);
|
|
if (!args.headed) {
|
|
process.stderr.write("Running headless. If sign-in is required, re-run with --headed.\n");
|
|
}
|
|
try {
|
|
cache = await captureToken({ teamId: args.teamId, channelId: args.channelId, headed: args.headed });
|
|
} catch (e) {
|
|
process.stderr.write(`Token capture failed: ${e.message}\n`);
|
|
process.exit(3);
|
|
}
|
|
saveCache(cache);
|
|
}
|
|
|
|
try {
|
|
const flat = await fetchMessages(cache, args);
|
|
await resolveDisplayNames(cache, flat, { teamId: args.teamId, channelId: args.channelId });
|
|
const grouped = args.includeReplies === false
|
|
? flat.map((m) => ({ ...m, replies: [] }))
|
|
: groupReplies(flat);
|
|
const rootCount = grouped.length;
|
|
const replyCount = grouped.reduce((n, r) => n + (r.replies?.length || 0), 0);
|
|
process.stdout.write(JSON.stringify({
|
|
value: grouped,
|
|
capturedAt: cache.capturedAt,
|
|
endpoint: cache.endpoint,
|
|
count: rootCount,
|
|
replyCount,
|
|
}, null, 2));
|
|
process.stdout.write("\n");
|
|
} catch (e) {
|
|
process.stderr.write(`Fetch failed: ${e.message}\n`);
|
|
process.exit(4);
|
|
}
|
|
})();
|
|
|