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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
34 changes: 34 additions & 0 deletions connectors/gmail/src/gmail.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -461,6 +461,40 @@ export class Gmail extends Connector<Gmail> {
}

private async setupMailboxWebhook(): Promise<void> {
// Tear down any prior watch and topic before creating new ones. Gmail
// enforces one watch per (mailbox, OAuth client) and returns 400
// "Only one user push notification client allowed per developer (call
// /stop then try again)" when users.watch() is called with a NEW topic
// while a watch is already active. setupMailboxWebhook always creates
// a fresh topic (createWebhook mints a new callback token → new topic
// name), so the existing watch must be stopped first; the orphaned
// Pub/Sub topic is also deleted to avoid leaking resources every
// self-heal renewal.
const existing = await this.get<MailboxWebhookState>("mailbox_webhook");
await this.clear("mailbox_webhook");
const cleanupApi = await this.getApiAny();
if (cleanupApi) {
try {
await cleanupApi.stopWatch();
} catch (error) {
// Best-effort — old watch may have already expired or never existed.
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: stopWatch (cleanup) failed`,
error
);
}
}
if (existing?.topicName) {
try {
await this.tools.network.deleteWebhook(existing.topicName);
} catch (error) {
console.warn(
`Gmail setupMailboxWebhook [${this.id}]: deleteWebhook (cleanup) failed`,
error
);
}
}

// createWebhook returns a Pub/Sub topic name when the provider is Google
// with Gmail scopes. The webhook delivers no extra args — onGmailWebhook
// operates on the single mailbox-wide watch.
Expand Down
174 changes: 105 additions & 69 deletions connectors/linkedin-messaging/src/linkedin.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -28,8 +28,16 @@ import {
type LinkedInProfile,
} from "@plotday/twister/tools/linkedin";

const CHANNEL_ID = "linkedin";
const CHANNEL_TITLE = "LinkedIn";
/**
* LinkedIn surfaces two distinct streams that users may want to enable
* independently: direct messages (high-volume, real conversations) and
* inbound connection requests (lower-volume, usually higher-noise). Each
* is its own Plot channel so users can opt into one without the other,
* and so the framework can track sync state, polling cadence, and
* read/unread separately.
*/
const CHANNEL_MESSAGES = "messages";
const CHANNEL_INVITATIONS = "invitations";

const TYPE_MESSAGE = "message";
const TYPE_INVITATION = "invitation";
Expand DownExpand Up@@ -98,7 +106,9 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {

readonly provider = AuthProvider.LinkedIn;
readonly scopes = LinkedInMessaging.SCOPES;
readonly singleChannel = true;
// Two channels: "messages" and "invitations". Users can enable either
// independently — connection-request triage and conversation triage are
// very different workflows and shouldn't share an on/off toggle.
readonly linkTypes = [
{
type: TYPE_MESSAGE,
Expand DownExpand Up@@ -139,14 +149,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
_auth: Authorization | null,
_token: AuthToken | null
): Promise<Channel[]> {
// Single implicit channel — the user's LinkedIn account is the channel.
// Splitting messages vs. invitations across multiple channels would
// expose a toggle that nobody actually wants to flip; keep it simple
// and bundle both into one stream.
return [{ id: CHANNEL_ID, title: CHANNEL_TITLE }];
return [
{ id: CHANNEL_MESSAGES, title: "Direct messages" },
{ id: CHANNEL_INVITATIONS, title: "Connection requests" },
];
}

async onChannelEnabled(channel: Channel): Promise<void> {
if (
channel.id !== CHANNEL_MESSAGES &&
channel.id !== CHANNEL_INVITATIONS
) {
return;
}
// Seed the sync state and queue the first batch as a task. Initial
// sync runs with the `initialSync` flag so the runtime suppresses
// notifications during the backfill.
Expand DownExpand Up@@ -182,14 +197,15 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

/**
* One adaptive-poll cycle. Runs in its own execution (queued via
* `runTask`) so each cycle gets a fresh ~1000-request budget.
* One adaptive-poll cycle for a channel. Runs in its own execution
* (queued via `runTask`) so each cycle gets a fresh ~1000-request
* budget.
*
* 1. List conversations updated since `lastSyncedActivityAt`.
* 2. For each conversation, list its new messages and save the link +
* notes.
* 3. Walk inbound connection invitations as a separate stream.
* 4. Schedule the next cycle based on activity recency.
* Branches on channel:
* - `messages`: list conversations updated since the cursor, fetch new
* messages for each, and save one link per conversation.
* - `invitations`: list inbound connection requests and save one link
* per invitation.
*/
async syncBatch(channelId: string): Promise<void> {
const state =
Expand All@@ -203,21 +219,59 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
? new Date(state.lastSyncedActivityAt)
: undefined;

let cursor: string | null = null;
let highWaterMark = state.lastSyncedActivityAt ?? 0;

if (channelId === CHANNEL_MESSAGES) {
highWaterMark = await this.syncMessages(state, since, highWaterMark);
} else if (channelId === CHANNEL_INVITATIONS) {
highWaterMark = await this.syncInvitations(state, highWaterMark);
} else {
// Unknown channel id — ignore so a stale runTask doesn't loop on it.
return;
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
});
}

/**
* Direct-messages sync. Walks conversations updated since the last
* `lastSyncedActivityAt` and saves one link per conversation, with the
* new messages attached as notes. Returns the updated high-water mark.
*/
private async syncMessages(
state: SyncState,
since: Date | undefined,
highWaterMark: number
): Promise<number> {
let cursor: string | null = null;
const conversationLinks: NewLinkWithNotes[] = [];

// Conversations
for (let page = 0; page < 5; page++) {
const result = await this.tools.linkedin.listConversations({
channelId,
channelId: CHANNEL_MESSAGES,
cursor,
since,
limit: 20,
});
for (const conv of result.conversations) {
const link = await this.buildConversationLink(
channelId,
conv,
state.initialSync,
since
Expand All@@ -234,39 +288,33 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
await this.tools.integrations.saveLinks(conversationLinks);
}

// Invitations (one page per cycle is plenty — most users see a few per
// day at most).
if (state.initialSync || sinceWasRecent(since)) {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(channelId, inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
}

// Persist progress.
await this.set(`sync_state_${channelId}`, {
initialSync: false,
lastSyncedActivityAt: highWaterMark || Date.now(),
lastUserActivityAt: state.lastUserActivityAt,
} satisfies SyncState);

// Tell the runtime we're done backfilling so the "syncing…" indicator
// clears on first pass.
if (state.initialSync) {
await this.tools.integrations.channelSyncCompleted(channelId);
}
return highWaterMark;
}

// Schedule the next cycle.
await this.scheduleNextSync(channelId, {
lastUserActivityAt: state.lastUserActivityAt,
lastSyncedActivityAt: highWaterMark,
/**
* Connection-requests sync. One page per cycle — most users see a few
* invitations per day at most, so a 20-item page is plenty. Each
* invitation becomes its own Plot link.
*/
private async syncInvitations(
state: SyncState,
highWaterMark: number
): Promise<number> {
const invitations = await this.tools.linkedin.listConnectionInvitations({
channelId: CHANNEL_INVITATIONS,
limit: 20,
});
const invitationLinks = invitations.invitations
.map((inv) => buildInvitationLink(inv, state.initialSync))
.filter((link): link is NewLinkWithNotes => link != null);
for (const inv of invitations.invitations) {
const ts = inv.sentAt.getTime();
if (ts > highWaterMark) highWaterMark = ts;
}
if (invitationLinks.length > 0) {
await this.tools.integrations.saveLinks(invitationLinks);
}
return highWaterMark;
}

private async scheduleNextSync(
Expand DownExpand Up@@ -300,15 +348,14 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
// ---------------------------------------------------------------------------

private async buildConversationLink(
channelId: string,
conv: LinkedInConversation,
initialSync: boolean,
since: Date | undefined
): Promise<NewLinkWithNotes | null> {
// Pull the new messages for this conversation. On initial sync we
// walk a single page (~20 messages) so the backfill stays bounded.
const messagesResult = await this.tools.linkedin.getMessages({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
since: initialSync ? undefined : since,
limit: 20,
Expand DownExpand Up@@ -340,7 +387,7 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn: conv.urn,
isGroup: conv.isGroup,
},
Expand All@@ -358,22 +405,19 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<NoteWriteBackResult | void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

const sent = await this.tools.linkedin.sendMessage({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
text: note.content ?? "",
});

// Bias the next poll cycle to be fresh — the user just replied so a
// response is likely incoming soon.
const state =
(await this.get<SyncState>(`sync_state_${channelId}`)) ?? null;
const state = await this.get<SyncState>(`sync_state_${CHANNEL_MESSAGES}`);
if (state) {
await this.set(`sync_state_${channelId}`, {
await this.set(`sync_state_${CHANNEL_MESSAGES}`, {
...state,
lastUserActivityAt: Date.now(),
} satisfies SyncState);
Expand All@@ -392,13 +436,11 @@ export class LinkedInMessaging extends Connector<LinkedInMessaging> {
): Promise<void> {
const meta = thread.meta ?? {};
const conversationUrn = meta.conversationUrn as string | undefined;
const channelId =
(meta.channelId as string | undefined) ?? CHANNEL_ID;
if (!conversationUrn) return;

try {
await this.tools.linkedin.markConversationRead({
channelId,
channelId: CHANNEL_MESSAGES,
conversationUrn,
read: !unread,
});
Expand DownExpand Up@@ -499,7 +541,6 @@ function joinParticipantNames(profiles: LinkedInProfile[]): string {
}

function buildInvitationLink(
channelId: string,
inv: LinkedInInvitation,
initialSync: boolean
): NewLinkWithNotes | null {
Expand DownExpand Up@@ -537,16 +578,11 @@ function buildInvitationLink(
notes,
meta: {
syncProvider: PROVIDER_KEY,
channelId,
channelId: CHANNEL_INVITATIONS,
invitationUrn: inv.urn,
sharedSecret: inv.sharedSecret,
inviterUrn: inv.inviter.urn,
},
...(initialSync ? { unread: false, archived: false } : {}),
} as NewLinkWithNotes;
}

function sinceWasRecent(since: Date | undefined): boolean {
if (!since) return true;
return Date.now() - since.getTime() < RECENT_WINDOW_MS;
}
Loading