mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-10-01 22:14:51 +02:00
[eric] mcp: X (Twitter) session-borrow MCP shim (15 tools, full human action set)
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
"""Module-level entrypoint so `python -m backend.apps.x_mcp_shim` works."""
|
||||
from backend.apps.x_mcp_shim.server import main
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,71 @@
|
||||
"""Dispatch each MCP tool call to the X client and format MCP content."""
|
||||
|
||||
import json
|
||||
from typing import Any, Dict
|
||||
|
||||
from backend.apps.social_shims.session_source import SessionUnavailable
|
||||
from backend.apps.x_mcp_shim import x_reads as reads
|
||||
from backend.apps.x_mcp_shim import x_writes as writes
|
||||
from backend.apps.x_mcp_shim.x_http import XError
|
||||
|
||||
|
||||
def mcp_ok(payload: Any) -> Dict[str, Any]:
|
||||
if isinstance(payload, str):
|
||||
return {"content": [{"type": "text", "text": payload}]}
|
||||
return {"content": [{"type": "text", "text": json.dumps(payload, indent=2, default=str)}]}
|
||||
|
||||
|
||||
def mcp_err(text: str) -> Dict[str, Any]:
|
||||
return {"content": [{"type": "text", "text": f"Error: {text}"}], "isError": True}
|
||||
|
||||
|
||||
def handle_tool_call(name: str, args: Dict[str, Any]) -> Dict[str, Any]:
|
||||
try:
|
||||
return mcp_ok(p_dispatch(name, args))
|
||||
except SessionUnavailable as e:
|
||||
return mcp_err(str(e))
|
||||
except XError as e:
|
||||
return mcp_err(str(e))
|
||||
except Exception as e:
|
||||
return mcp_err(f"x shim error: {e!r}")
|
||||
|
||||
|
||||
def p_dispatch(name: str, a: Dict[str, Any]) -> Any:
|
||||
if name == "x_whoami":
|
||||
return reads.whoami()
|
||||
if name == "x_timeline":
|
||||
return reads.timeline(a.get("kind", "foryou"), p_lim(a.get("count"), 20), a.get("cursor", ""))
|
||||
if name == "x_user_tweets":
|
||||
return reads.user_tweets(a.get("username", ""), p_lim(a.get("count"), 20), a.get("cursor", ""))
|
||||
if name == "x_get_tweet":
|
||||
return reads.get_tweet(a.get("target", ""), p_lim(a.get("replies_limit"), 30))
|
||||
if name == "x_search":
|
||||
return reads.search(a.get("query", ""), a.get("product", "top"), p_lim(a.get("count"), 20), a.get("cursor", ""))
|
||||
if name == "x_get_user":
|
||||
return reads.get_user(a.get("username", ""))
|
||||
if name == "x_bookmarks":
|
||||
return reads.bookmarks(p_lim(a.get("count"), 20), a.get("cursor", ""))
|
||||
if name == "x_notifications":
|
||||
return reads.notifications(p_lim(a.get("count"), 20), a.get("cursor", ""))
|
||||
if name == "x_tweet":
|
||||
return writes.tweet(a.get("text", ""), a.get("reply_to", ""), a.get("quote_id", ""))
|
||||
if name == "x_delete_tweet":
|
||||
return writes.delete_tweet(a.get("target", ""))
|
||||
if name == "x_like":
|
||||
return writes.like(a.get("target", ""), bool(a.get("unlike")))
|
||||
if name == "x_retweet":
|
||||
return writes.retweet(a.get("target", ""), bool(a.get("undo")))
|
||||
if name == "x_bookmark":
|
||||
return writes.bookmark(a.get("target", ""), bool(a.get("remove")))
|
||||
if name == "x_follow":
|
||||
return writes.follow(a.get("username", ""), bool(a.get("unfollow")))
|
||||
if name == "x_send_dm":
|
||||
return writes.send_dm(a.get("recipient", ""), a.get("text", ""))
|
||||
raise XError(f"Unknown tool: {name}")
|
||||
|
||||
|
||||
def p_lim(v: Any, default: int) -> int:
|
||||
try:
|
||||
return max(1, min(int(v), 100))
|
||||
except (TypeError, ValueError):
|
||||
return default
|
||||
@@ -0,0 +1,33 @@
|
||||
"""X's per-action pacing config on top of the shared RateLimiter.
|
||||
|
||||
Reads are generous; posting is deliberately slow, likes/retweets moderate, follows
|
||||
and DMs slow, so the account never bursts like a bot. The shared core owns the
|
||||
algorithm + honors X's x-rate-limit-* headers and 429 backoff.
|
||||
"""
|
||||
|
||||
from typing import Dict, Tuple
|
||||
|
||||
from backend.apps.social_shims.rate_limit_core import RateLimiter
|
||||
|
||||
# action -> (bucket_capacity, seconds_to_refill_one_token).
|
||||
BUCKETS: Dict[str, Tuple[float, float]] = {
|
||||
"read": (30.0, 1.0),
|
||||
"tweet": (3.0, 60.0),
|
||||
"like": (20.0, 2.0),
|
||||
"follow": (8.0, 6.0),
|
||||
"dm": (5.0, 20.0),
|
||||
}
|
||||
|
||||
p_limiter = RateLimiter(BUCKETS, min_gap_s=1.0, jitter_s=0.7)
|
||||
|
||||
|
||||
def acquire(action: str) -> None:
|
||||
p_limiter.acquire(action)
|
||||
|
||||
|
||||
def note_response(status: int, headers: Dict[str, str]) -> None:
|
||||
p_limiter.note_response(status, headers)
|
||||
|
||||
|
||||
def bucket_for(action: str) -> str:
|
||||
return p_limiter.bucket_for(action)
|
||||
@@ -0,0 +1,63 @@
|
||||
"""Stdio JSON-RPC MCP server for X (Twitter).
|
||||
|
||||
Mirrors the discord/reddit shim loop: stdlib-only, no backend.config imports, so the
|
||||
subprocess starts fast. The tool surface lives in tools.py; dispatch in handlers.py.
|
||||
"""
|
||||
|
||||
import json
|
||||
import sys
|
||||
from typing import Any, Optional
|
||||
|
||||
from backend.apps.x_mcp_shim.handlers import handle_tool_call, mcp_err
|
||||
from backend.apps.x_mcp_shim.tools import TOOLS
|
||||
|
||||
|
||||
def p_send(id_: Any, result: Optional[dict] = None, error: Optional[dict] = None) -> None:
|
||||
msg: dict = {"jsonrpc": "2.0", "id": id_}
|
||||
if error is not None:
|
||||
msg["error"] = error
|
||||
else:
|
||||
msg["result"] = result
|
||||
sys.stdout.write(json.dumps(msg) + "\n")
|
||||
sys.stdout.flush()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
for line in sys.stdin:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
msg = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
|
||||
method = msg.get("method")
|
||||
id_ = msg.get("id")
|
||||
params = msg.get("params", {}) or {}
|
||||
|
||||
if method == "initialize":
|
||||
p_send(id_, {
|
||||
"protocolVersion": "2024-11-05",
|
||||
"capabilities": {"tools": {}},
|
||||
"serverInfo": {"name": "openswarm-x", "version": "1.0.0"},
|
||||
})
|
||||
elif method == "notifications/initialized":
|
||||
pass
|
||||
elif method == "tools/list":
|
||||
p_send(id_, {"tools": TOOLS})
|
||||
elif method == "tools/call":
|
||||
name = params.get("name", "")
|
||||
args = params.get("arguments", {}) or {}
|
||||
try:
|
||||
p_send(id_, handle_tool_call(name, args))
|
||||
except Exception as e:
|
||||
p_send(id_, mcp_err(f"shim crashed: {e!r}"))
|
||||
elif method == "ping":
|
||||
p_send(id_, {})
|
||||
elif id_ is not None:
|
||||
p_send(id_, error={"code": -32601, "message": f"Method not found: {method}"})
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,180 @@
|
||||
"""MCP tool surface for X (Twitter): the full set of things a logged-in human can do.
|
||||
|
||||
Reads (timeline/search/tweet/user/bookmarks/notifications) and writes (tweet with
|
||||
reply+quote, delete, like, retweet, bookmark, follow, DM). A tweet `target` may be a
|
||||
status URL, a numeric id, or a t-prefixed id; the read tools return ids you pass back.
|
||||
"""
|
||||
|
||||
OBJ = "object"
|
||||
|
||||
TOOLS = [
|
||||
{
|
||||
"name": "x_whoami",
|
||||
"description": "Confirm the logged-in X session and return your handle + profile. Use first to verify the session is live.",
|
||||
"inputSchema": {"type": OBJ, "properties": {}},
|
||||
},
|
||||
{
|
||||
"name": "x_timeline",
|
||||
"description": "Read your home timeline (the For You or Following feed).",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"kind": {"type": "string", "enum": ["foryou", "following"], "default": "foryou"},
|
||||
"count": {"type": "integer", "default": 20, "description": "1-100"},
|
||||
"cursor": {"type": "string", "description": "Pagination cursor from a previous call."},
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_user_tweets",
|
||||
"description": "Read a user's recent tweets and replies by @handle.",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"username": {"type": "string", "description": "Handle, with or without @."},
|
||||
"count": {"type": "integer", "default": 20},
|
||||
"cursor": {"type": "string"},
|
||||
},
|
||||
"required": ["username"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_get_tweet",
|
||||
"description": "Get a tweet plus its replies. target is a status URL, a numeric id, or a t-prefixed id.",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"target": {"type": "string"},
|
||||
"replies_limit": {"type": "integer", "default": 30},
|
||||
},
|
||||
"required": ["target"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_search",
|
||||
"description": "Search tweets. product: top, latest, people, or media.",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"query": {"type": "string"},
|
||||
"product": {"type": "string", "enum": ["top", "latest", "people", "media"], "default": "top"},
|
||||
"count": {"type": "integer", "default": 20},
|
||||
"cursor": {"type": "string"},
|
||||
},
|
||||
"required": ["query"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_get_user",
|
||||
"description": "Get a user's profile (bio, follower/following/tweet counts) by @handle.",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {"username": {"type": "string"}},
|
||||
"required": ["username"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_bookmarks",
|
||||
"description": "List your bookmarked tweets.",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"count": {"type": "integer", "default": 20},
|
||||
"cursor": {"type": "string"},
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_notifications",
|
||||
"description": "Read your notifications (mentions, likes, follows, replies).",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"count": {"type": "integer", "default": 20},
|
||||
"cursor": {"type": "string"},
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_tweet",
|
||||
"description": "Post a tweet. Set reply_to to reply to a tweet, or quote_id to quote one (both accept a URL/id).",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"text": {"type": "string"},
|
||||
"reply_to": {"type": "string", "description": "Tweet URL/id to reply to."},
|
||||
"quote_id": {"type": "string", "description": "Tweet URL/id to quote."},
|
||||
},
|
||||
"required": ["text"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_delete_tweet",
|
||||
"description": "Delete one of your own tweets by URL/id.",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {"target": {"type": "string"}},
|
||||
"required": ["target"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_like",
|
||||
"description": "Like a tweet (or unlike with unlike=true).",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"target": {"type": "string"},
|
||||
"unlike": {"type": "boolean", "default": False},
|
||||
},
|
||||
"required": ["target"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_retweet",
|
||||
"description": "Retweet a tweet (or undo with undo=true).",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"target": {"type": "string"},
|
||||
"undo": {"type": "boolean", "default": False},
|
||||
},
|
||||
"required": ["target"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_bookmark",
|
||||
"description": "Bookmark a tweet (or remove with remove=true).",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"target": {"type": "string"},
|
||||
"remove": {"type": "boolean", "default": False},
|
||||
},
|
||||
"required": ["target"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_follow",
|
||||
"description": "Follow a user by @handle (or unfollow with unfollow=true).",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"username": {"type": "string"},
|
||||
"unfollow": {"type": "boolean", "default": False},
|
||||
},
|
||||
"required": ["username"],
|
||||
},
|
||||
},
|
||||
{
|
||||
"name": "x_send_dm",
|
||||
"description": "Send a direct message. recipient is a @handle or a numeric user id.",
|
||||
"inputSchema": {
|
||||
"type": OBJ,
|
||||
"properties": {
|
||||
"recipient": {"type": "string"},
|
||||
"text": {"type": "string"},
|
||||
},
|
||||
"required": ["recipient", "text"],
|
||||
},
|
||||
},
|
||||
]
|
||||
@@ -0,0 +1,75 @@
|
||||
"""X (Twitter) web-client constants: the public bearer + the GraphQL operation map.
|
||||
|
||||
The Authorization bearer below is the PUBLIC token x.com ships to every web client
|
||||
(logged-in or not); it is not a secret and not per-user. Real auth is the borrowed
|
||||
auth_token + ct0 cookies. The GraphQL queryIds drift whenever X redeploys its web
|
||||
app: refresh them by opening x.com in the OpenSwarm browser, watching the Network
|
||||
tab for /i/api/graphql/<id>/<OpName>, and pasting the new <id> here. This is the one
|
||||
drift-prone surface, deliberately isolated so a refresh is a one-line edit.
|
||||
|
||||
Known gap: X's newest anti-automation header (x-client-transaction-id) is generated
|
||||
by obfuscated client JS we don't replicate; some endpoints may 404/403 without it.
|
||||
That's the X equivalent of Reddit's bearer-harvest assumption: structurally sound,
|
||||
not live-proven here.
|
||||
"""
|
||||
|
||||
# Public web-app bearer (constant across all x.com web clients; not a credential).
|
||||
WEB_BEARER = (
|
||||
"AAAAAAAAAAAAAAAAAAAAANRILgAAAAAAnNwIzUejRCOuH5E6I8xnZz4puTs%3D"
|
||||
"1Zv7ttfk8LF81IUq16cHjhLTvJu4FA33AGWWjCpTnA"
|
||||
)
|
||||
|
||||
# OpName -> queryId. Drift-prone; see module docstring to refresh from a live capture.
|
||||
GRAPHQL_IDS = {
|
||||
"UserByScreenName": "G3KGOASz96M-Qu0nwmGXNg",
|
||||
"UserTweets": "E3opETHurmVJflFsUBVuUQ",
|
||||
"TweetDetail": "xOhkmRac04YFZmOzU9PJHg",
|
||||
"SearchTimeline": "nKAncKPF1fV1xltvF3UUlw",
|
||||
"HomeTimeline": "uPv755D929tshj6KsxkSZg",
|
||||
"HomeLatestTimeline": "8Rfm0g9b2-9La8Rmd1IPzw",
|
||||
"Bookmarks": "j5KExFXxK0Nz1tQNXEx6KQ",
|
||||
"CreateTweet": "znq5dRMnAYIRgIBQhGCRkg",
|
||||
"DeleteTweet": "VaenaVgh5q5ih7kvyVjgtg",
|
||||
"FavoriteTweet": "lI07N6Otwv1PhnEgXILM7A",
|
||||
"UnfavoriteTweet": "ZYKSe-w7KEslx3JhSIk5LA",
|
||||
"CreateRetweet": "ojPdsZsimiJrUGLR1sjUtA",
|
||||
"DeleteRetweet": "iQtK4dl5hBmXewYZuEOKVw",
|
||||
"CreateBookmark": "aoDbu3RHznuiSkQ9aNM67Q",
|
||||
"DeleteBookmark": "Wlmlj2-xzyS1GN3a6cj-mQ",
|
||||
}
|
||||
|
||||
# Feature flags X's GraphQL requires; a missing key 400s with "features cannot be null".
|
||||
# Also drift-prone; kept broad. Refresh alongside the queryIds.
|
||||
DEFAULT_FEATURES = {
|
||||
"rweb_video_screen_enabled": False,
|
||||
"profile_label_improvements_pcf_label_in_post_enabled": True,
|
||||
"responsive_web_graphql_exclude_directive_enabled": True,
|
||||
"verified_phone_label_enabled": False,
|
||||
"creator_subscriptions_tweet_preview_api_enabled": True,
|
||||
"responsive_web_graphql_timeline_navigation_enabled": True,
|
||||
"responsive_web_graphql_skip_user_profile_image_extensions_enabled": False,
|
||||
"premium_content_api_read_enabled": False,
|
||||
"communities_web_enable_tweet_community_results_fetch": True,
|
||||
"c9s_tweet_anatomy_moderator_badge_enabled": True,
|
||||
"responsive_web_grok_analyze_button_fetch_trends_enabled": False,
|
||||
"responsive_web_grok_analyze_post_followups_enabled": True,
|
||||
"responsive_web_jetfuel_frame": False,
|
||||
"responsive_web_grok_share_attachment_enabled": True,
|
||||
"articles_preview_enabled": True,
|
||||
"responsive_web_edit_tweet_api_enabled": True,
|
||||
"graphql_is_translatable_rweb_tweet_is_translatable_enabled": True,
|
||||
"view_counts_everywhere_api_enabled": True,
|
||||
"longform_notetweets_consumption_enabled": True,
|
||||
"responsive_web_twitter_article_tweet_consumption_enabled": True,
|
||||
"tweet_awards_web_tipping_enabled": False,
|
||||
"responsive_web_grok_show_grok_translated_post": False,
|
||||
"responsive_web_grok_analysis_button_from_backend": True,
|
||||
"creator_subscriptions_quote_tweet_preview_enabled": False,
|
||||
"freedom_of_speech_not_reach_fetch_enabled": True,
|
||||
"standardized_nudges_misinfo": True,
|
||||
"tweet_with_visibility_results_prefer_gql_limited_actions_policy_enabled": True,
|
||||
"longform_notetweets_rich_text_read_enabled": True,
|
||||
"longform_notetweets_inline_media_enabled": True,
|
||||
"responsive_web_grok_image_annotation_enabled": True,
|
||||
"responsive_web_enhance_cards_enabled": False,
|
||||
}
|
||||
@@ -0,0 +1,104 @@
|
||||
"""Low-level authed X (Twitter) transport.
|
||||
|
||||
Borrow the user's x.com session (auth_token + ct0 cookies), attach the public web
|
||||
bearer + the ct0-derived CSRF header, and call x.com's own /i/api GraphQL + v1.1/v2
|
||||
surfaces exactly as the logged-in web client does. Rate-limited and self-refreshing
|
||||
on a 401/403 by re-borrowing the session. stdlib-only to match the sibling shims.
|
||||
"""
|
||||
|
||||
import json
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
from backend.apps.social_shims.session_source import cookie_value, get_session, invalidate
|
||||
from backend.apps.x_mcp_shim import rate_limit
|
||||
from backend.apps.x_mcp_shim.x_endpoints import DEFAULT_FEATURES, GRAPHQL_IDS, WEB_BEARER
|
||||
|
||||
DOMAIN = "x.com"
|
||||
API = "https://x.com/i/api"
|
||||
GRAPHQL = f"{API}/graphql"
|
||||
|
||||
|
||||
class XError(Exception):
|
||||
"""An X request failed in a way worth surfacing to the agent."""
|
||||
|
||||
|
||||
def p_send(method: str, url: str, *, data: Optional[bytes], content_type: Optional[str],
|
||||
action: str, retried: bool = False) -> Any:
|
||||
rate_limit.acquire(action)
|
||||
cookie, ua = get_session(DOMAIN)
|
||||
ct0 = cookie_value(DOMAIN, "ct0")
|
||||
if not ct0:
|
||||
invalidate(DOMAIN)
|
||||
raise XError("No x.com CSRF cookie (ct0). Open x.com in the OpenSwarm browser, sign in, then retry.")
|
||||
headers = {
|
||||
"Authorization": f"Bearer {WEB_BEARER}",
|
||||
"Cookie": cookie,
|
||||
"User-Agent": ua,
|
||||
"x-csrf-token": ct0,
|
||||
"x-twitter-auth-type": "OAuth2Session",
|
||||
"x-twitter-active-user": "yes",
|
||||
"x-twitter-client-language": "en",
|
||||
"Accept": "application/json",
|
||||
"Referer": "https://x.com/",
|
||||
}
|
||||
if data is not None and content_type:
|
||||
headers["Content-Type"] = content_type
|
||||
req = urllib.request.Request(url, data=data, headers=headers, method=method)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=30.0) as resp:
|
||||
status, raw, rhdr = resp.status, resp.read(), dict(resp.headers)
|
||||
except urllib.error.HTTPError as e:
|
||||
status, raw, rhdr = e.code, (e.read() if e.fp else b""), dict(e.headers or {})
|
||||
except urllib.error.URLError as e:
|
||||
raise XError(f"x.com unreachable: {getattr(e, 'reason', e)}")
|
||||
|
||||
rate_limit.note_response(status, {k.lower(): v for k, v in rhdr.items()})
|
||||
if status in (401, 403) and not retried:
|
||||
invalidate(DOMAIN)
|
||||
return p_send(method, url, data=data, content_type=content_type, action=action, retried=True)
|
||||
if status == 429:
|
||||
raise XError("x.com is rate-limiting this account; slow down and retry shortly.")
|
||||
if status >= 400:
|
||||
raise XError(f"x.com HTTP {status}: {raw[:300].decode('utf-8', 'replace')}")
|
||||
try:
|
||||
return json.loads(raw.decode("utf-8", errors="replace") or "{}")
|
||||
except json.JSONDecodeError:
|
||||
return {"raw": raw.decode("utf-8", errors="replace")}
|
||||
|
||||
|
||||
def graphql(op: str, variables: Dict[str, Any], *, method: str = "GET",
|
||||
features: bool = True, action: str = "read") -> Any:
|
||||
"""Call a GraphQL operation by name, looking up its (drift-prone) queryId."""
|
||||
qid = GRAPHQL_IDS.get(op)
|
||||
if not qid:
|
||||
raise XError(f"Unknown GraphQL op {op!r}; add its queryId to x_endpoints.GRAPHQL_IDS.")
|
||||
url = f"{GRAPHQL}/{qid}/{op}"
|
||||
if method == "GET":
|
||||
params = {"variables": json.dumps(variables, separators=(",", ":"))}
|
||||
if features:
|
||||
params["features"] = json.dumps(DEFAULT_FEATURES, separators=(",", ":"))
|
||||
return p_send("GET", url + "?" + urllib.parse.urlencode(params),
|
||||
data=None, content_type=None, action=action)
|
||||
body: Dict[str, Any] = {"variables": variables, "queryId": qid}
|
||||
if features:
|
||||
body["features"] = DEFAULT_FEATURES
|
||||
return p_send("POST", url, data=json.dumps(body).encode(), content_type="application/json", action=action)
|
||||
|
||||
|
||||
def rest(method: str, path: str, *, params: Optional[Dict[str, Any]] = None,
|
||||
form: Optional[Dict[str, Any]] = None, json_body: Optional[Dict[str, Any]] = None,
|
||||
action: str = "read") -> Any:
|
||||
"""Call a legacy v1.1/v2 endpoint (more stable than GraphQL for follow/DM). path includes the version."""
|
||||
url = f"{API}/{path}"
|
||||
if params:
|
||||
url += "?" + urllib.parse.urlencode({k: v for k, v in params.items() if v is not None})
|
||||
if json_body is not None:
|
||||
return p_send(method, url, data=json.dumps(json_body).encode(),
|
||||
content_type="application/json", action=action)
|
||||
if form is not None:
|
||||
data = urllib.parse.urlencode({k: v for k, v in form.items() if v is not None}).encode()
|
||||
return p_send(method, url, data=data, content_type="application/x-www-form-urlencoded", action=action)
|
||||
return p_send(method, url, data=None, content_type=None, action=action)
|
||||
@@ -0,0 +1,191 @@
|
||||
"""Read operations over x.com's own /i/api GraphQL + v1.1/v2 surfaces.
|
||||
|
||||
Returns compact, token-frugal tweet/user records (truncated text) instead of X's
|
||||
deeply-nested GraphQL firehose, so the agent sees what a human skims. The parsing
|
||||
walks for `tweet_results` nodes anywhere in the tree, which survives X's frequent
|
||||
timeline-shape reshuffles better than fixed index paths.
|
||||
"""
|
||||
|
||||
import re
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from backend.apps.x_mcp_shim.x_http import XError, graphql, rest
|
||||
|
||||
TEXT_CAP = 1200
|
||||
|
||||
|
||||
def p_trunc(s: Optional[str]) -> str:
|
||||
s = s or ""
|
||||
return s if len(s) <= TEXT_CAP else s[:TEXT_CAP] + f"... [+{len(s) - TEXT_CAP} chars]"
|
||||
|
||||
|
||||
def tweet_id_of(target: str) -> str:
|
||||
"""Accept a status URL, a t-prefixed id, or a bare id; return the numeric id."""
|
||||
m = re.search(r"(\d{5,})", target or "")
|
||||
return m.group(1) if m else (target or "")
|
||||
|
||||
|
||||
def normalize_tweet(result: Any) -> Optional[Dict[str, Any]]:
|
||||
if not isinstance(result, dict):
|
||||
return None
|
||||
if result.get("__typename") == "TweetWithVisibilityResults":
|
||||
result = result.get("tweet", result)
|
||||
legacy = result.get("legacy") or {}
|
||||
if not legacy and not result.get("rest_id"):
|
||||
return None
|
||||
user_result = ((result.get("core") or {}).get("user_results") or {}).get("result") or {}
|
||||
user_legacy = user_result.get("legacy") or {}
|
||||
user_core = user_result.get("core") or {}
|
||||
note = ((result.get("note_tweet") or {}).get("note_tweet_results") or {}).get("result") or {}
|
||||
text = note.get("text") or legacy.get("full_text") or ""
|
||||
return {
|
||||
"id": result.get("rest_id") or legacy.get("id_str"),
|
||||
"author": user_legacy.get("screen_name") or user_core.get("screen_name"),
|
||||
"name": user_legacy.get("name") or user_core.get("name"),
|
||||
"text": p_trunc(text),
|
||||
"likes": legacy.get("favorite_count"),
|
||||
"retweets": legacy.get("retweet_count"),
|
||||
"replies": legacy.get("reply_count"),
|
||||
"quotes": legacy.get("quote_count"),
|
||||
"views": (result.get("views") or {}).get("count"),
|
||||
"created_at": legacy.get("created_at"),
|
||||
"lang": legacy.get("lang"),
|
||||
}
|
||||
|
||||
|
||||
def p_collect_tweets(node: Any, out: List[Dict[str, Any]], cap: int) -> None:
|
||||
if len(out) >= cap:
|
||||
return
|
||||
if isinstance(node, dict):
|
||||
tr = node.get("tweet_results")
|
||||
if isinstance(tr, dict) and isinstance(tr.get("result"), dict):
|
||||
t = normalize_tweet(tr["result"])
|
||||
if t and t.get("id") and not any(x["id"] == t["id"] for x in out):
|
||||
out.append(t)
|
||||
for v in node.values():
|
||||
p_collect_tweets(v, out, cap)
|
||||
elif isinstance(node, list):
|
||||
for v in node:
|
||||
p_collect_tweets(v, out, cap)
|
||||
|
||||
|
||||
def p_cursor(node: Any) -> Optional[str]:
|
||||
found: List[str] = []
|
||||
|
||||
def walk(n: Any) -> None:
|
||||
if found:
|
||||
return
|
||||
if isinstance(n, dict):
|
||||
if n.get("cursorType") == "Bottom" and n.get("value"):
|
||||
found.append(n["value"])
|
||||
for v in n.values():
|
||||
walk(v)
|
||||
elif isinstance(n, list):
|
||||
for v in n:
|
||||
walk(v)
|
||||
|
||||
walk(node)
|
||||
return found[0] if found else None
|
||||
|
||||
|
||||
def p_timeline_out(resp: Any, cap: int) -> Dict[str, Any]:
|
||||
items: List[Dict[str, Any]] = []
|
||||
p_collect_tweets(resp, items, cap)
|
||||
return {"tweets": items, "cursor": p_cursor(resp)}
|
||||
|
||||
|
||||
def get_user(screen_name: str) -> Dict[str, Any]:
|
||||
resp = graphql("UserByScreenName", {"screen_name": screen_name.lstrip("@")})
|
||||
result = (((resp or {}).get("data") or {}).get("user") or {}).get("result") or {}
|
||||
legacy = result.get("legacy") or {}
|
||||
core = result.get("core") or {}
|
||||
return {
|
||||
"id": result.get("rest_id"),
|
||||
"screen_name": legacy.get("screen_name") or core.get("screen_name"),
|
||||
"name": legacy.get("name") or core.get("name"),
|
||||
"bio": p_trunc(legacy.get("description")),
|
||||
"followers": legacy.get("followers_count"),
|
||||
"following": legacy.get("friends_count"),
|
||||
"tweets": legacy.get("statuses_count"),
|
||||
"verified": result.get("is_blue_verified") or legacy.get("verified"),
|
||||
"created_at": legacy.get("created_at") or core.get("created_at"),
|
||||
}
|
||||
|
||||
|
||||
def resolve_user_id(screen_name: str) -> str:
|
||||
uid = get_user(screen_name).get("id")
|
||||
if not uid:
|
||||
raise XError(f"Could not resolve @{screen_name.lstrip('@')} to a user id.")
|
||||
return str(uid)
|
||||
|
||||
|
||||
def whoami() -> Dict[str, Any]:
|
||||
settings = rest("GET", "1.1/account/settings.json")
|
||||
screen = settings.get("screen_name", "")
|
||||
out: Dict[str, Any] = {"screen_name": screen}
|
||||
if screen:
|
||||
try:
|
||||
out["profile"] = get_user(screen)
|
||||
except XError:
|
||||
pass
|
||||
return out
|
||||
|
||||
|
||||
def timeline(kind: str, count: int, cursor: str) -> Dict[str, Any]:
|
||||
op = "HomeLatestTimeline" if kind == "following" else "HomeTimeline"
|
||||
variables: Dict[str, Any] = {"count": count, "includePromotedContent": False,
|
||||
"latestControlAvailable": True, "withCommunity": True}
|
||||
if cursor:
|
||||
variables["cursor"] = cursor
|
||||
return p_timeline_out(graphql(op, variables, method="POST"), count)
|
||||
|
||||
|
||||
def user_tweets(screen_name: str, count: int, cursor: str) -> Dict[str, Any]:
|
||||
uid = resolve_user_id(screen_name)
|
||||
variables: Dict[str, Any] = {"userId": uid, "count": count, "includePromotedContent": False,
|
||||
"withQuickPromoteEligibilityTweetFields": False, "withVoice": True,
|
||||
"withV2Timeline": True}
|
||||
if cursor:
|
||||
variables["cursor"] = cursor
|
||||
return p_timeline_out(graphql("UserTweets", variables), count)
|
||||
|
||||
|
||||
def get_tweet(target: str, count: int) -> Dict[str, Any]:
|
||||
focal = tweet_id_of(target)
|
||||
variables: Dict[str, Any] = {"focalTweetId": focal, "with_rux_injections": False,
|
||||
"includePromotedContent": False, "withCommunity": True,
|
||||
"withQuickPromoteEligibilityTweetFields": True, "withBirdwatchNotes": True,
|
||||
"withVoice": True, "withV2Timeline": True}
|
||||
out = p_timeline_out(graphql("TweetDetail", variables), count + 1)
|
||||
tweets = out["tweets"]
|
||||
main = next((t for t in tweets if t["id"] == focal), tweets[0] if tweets else {})
|
||||
replies = [t for t in tweets if t.get("id") != main.get("id")]
|
||||
return {"tweet": main, "replies": replies[:count]}
|
||||
|
||||
|
||||
def search(query: str, product: str, count: int, cursor: str) -> Dict[str, Any]:
|
||||
product = product.capitalize() if (product or "").lower() in ("top", "latest", "people", "media") else "Top"
|
||||
variables: Dict[str, Any] = {"rawQuery": query, "count": count, "querySource": "typed_query", "product": product}
|
||||
if cursor:
|
||||
variables["cursor"] = cursor
|
||||
return p_timeline_out(graphql("SearchTimeline", variables), count)
|
||||
|
||||
|
||||
def bookmarks(count: int, cursor: str) -> Dict[str, Any]:
|
||||
variables: Dict[str, Any] = {"count": count, "includePromotedContent": False}
|
||||
if cursor:
|
||||
variables["cursor"] = cursor
|
||||
return p_timeline_out(graphql("Bookmarks", variables), count)
|
||||
|
||||
|
||||
def notifications(count: int, cursor: str) -> Dict[str, Any]:
|
||||
resp = rest("GET", "2/notifications/all.json", params={"count": count, "cursor": cursor or None})
|
||||
notes = (resp or {}).get("globalObjects", {}).get("notifications", {})
|
||||
out = []
|
||||
for nid, n in list(notes.items())[:count]:
|
||||
out.append({
|
||||
"id": nid,
|
||||
"text": (n.get("message") or {}).get("text"),
|
||||
"timestamp_ms": n.get("timestampMs"),
|
||||
})
|
||||
return {"notifications": out}
|
||||
@@ -0,0 +1,74 @@
|
||||
"""Write operations: everything a logged-in human does on X.
|
||||
|
||||
Tweets (with reply/quote), deletes, likes, retweets, bookmarks, follows, and DMs,
|
||||
all via the user's borrowed session. Each call rides the rate limiter's write buckets.
|
||||
"""
|
||||
|
||||
from typing import Any, Dict
|
||||
|
||||
from backend.apps.x_mcp_shim.x_http import graphql, rest
|
||||
from backend.apps.x_mcp_shim.x_reads import normalize_tweet, resolve_user_id, tweet_id_of
|
||||
|
||||
|
||||
def p_created_tweet(resp: Any) -> Dict[str, Any]:
|
||||
result = ((((resp or {}).get("data") or {}).get("create_tweet") or {})
|
||||
.get("tweet_results") or {}).get("result") or {}
|
||||
t = normalize_tweet(result) or {}
|
||||
return {"id": t.get("id"), "text": t.get("text")}
|
||||
|
||||
|
||||
def tweet(text: str, reply_to: str, quote_id: str) -> Dict[str, Any]:
|
||||
variables: Dict[str, Any] = {
|
||||
"tweet_text": text,
|
||||
"dark_request": False,
|
||||
"media": {"media_entities": [], "possibly_sensitive": False},
|
||||
"semantic_annotation_ids": [],
|
||||
}
|
||||
if reply_to:
|
||||
variables["reply"] = {"in_reply_to_tweet_id": tweet_id_of(reply_to), "exclude_reply_user_ids": []}
|
||||
if quote_id:
|
||||
variables["attachment_url"] = f"https://x.com/i/status/{tweet_id_of(quote_id)}"
|
||||
return p_created_tweet(graphql("CreateTweet", variables, method="POST", action="tweet"))
|
||||
|
||||
|
||||
def delete_tweet(target: str) -> Dict[str, Any]:
|
||||
tid = tweet_id_of(target)
|
||||
graphql("DeleteTweet", {"tweet_id": tid, "dark_request": False}, method="POST", action="tweet")
|
||||
return {"id": tid, "deleted": True}
|
||||
|
||||
|
||||
def like(target: str, unlike: bool) -> Dict[str, Any]:
|
||||
tid = tweet_id_of(target)
|
||||
graphql("UnfavoriteTweet" if unlike else "FavoriteTweet", {"tweet_id": tid}, method="POST", action="like")
|
||||
return {"id": tid, "liked": not unlike}
|
||||
|
||||
|
||||
def retweet(target: str, undo: bool) -> Dict[str, Any]:
|
||||
tid = tweet_id_of(target)
|
||||
if undo:
|
||||
graphql("DeleteRetweet", {"source_tweet_id": tid, "dark_request": False}, method="POST", action="like")
|
||||
else:
|
||||
graphql("CreateRetweet", {"tweet_id": tid, "dark_request": False}, method="POST", action="like")
|
||||
return {"id": tid, "retweeted": not undo}
|
||||
|
||||
|
||||
def bookmark(target: str, remove: bool) -> Dict[str, Any]:
|
||||
tid = tweet_id_of(target)
|
||||
graphql("DeleteBookmark" if remove else "CreateBookmark", {"tweet_id": tid}, method="POST", action="like")
|
||||
return {"id": tid, "bookmarked": not remove}
|
||||
|
||||
|
||||
def follow(screen_name: str, unfollow: bool) -> Dict[str, Any]:
|
||||
uid = resolve_user_id(screen_name)
|
||||
path = "1.1/friendships/destroy.json" if unfollow else "1.1/friendships/create.json"
|
||||
rest("POST", path, form={"user_id": uid}, action="follow")
|
||||
return {"screen_name": screen_name.lstrip("@"), "following": not unfollow}
|
||||
|
||||
|
||||
def send_dm(recipient: str, text: str) -> Dict[str, Any]:
|
||||
rid = recipient if recipient.isdigit() else resolve_user_id(recipient)
|
||||
body = {"event": {"type": "message_create",
|
||||
"message_create": {"target": {"recipient_id": rid},
|
||||
"message_data": {"text": text}}}}
|
||||
rest("POST", "1.1/dm/new2.json", json_body=body, action="dm")
|
||||
return {"to": recipient, "sent": True}
|
||||
@@ -0,0 +1,138 @@
|
||||
"""Unit coverage for the X (Twitter) MCP shim's pure logic: tweet-id extraction, the
|
||||
deeply-nested GraphQL tweet/cursor walker, the rate limiter, and tool dispatch with
|
||||
the network mocked. Live posting can't be verified without a logged-in session, so the
|
||||
GraphQL contract is pinned here against canned x.com payload shapes."""
|
||||
|
||||
import json
|
||||
import time
|
||||
|
||||
from unittest.mock import patch
|
||||
|
||||
from backend.apps.social_shims.session_source import SessionUnavailable
|
||||
from backend.apps.x_mcp_shim import rate_limit, x_reads, x_writes
|
||||
from backend.apps.x_mcp_shim.handlers import handle_tool_call
|
||||
from backend.apps.x_mcp_shim.x_reads import normalize_tweet, tweet_id_of
|
||||
|
||||
|
||||
def p_text(result: dict) -> str:
|
||||
return result["content"][0]["text"]
|
||||
|
||||
|
||||
CANNED_TWEET = {
|
||||
"__typename": "Tweet",
|
||||
"rest_id": "111",
|
||||
"legacy": {"full_text": "hello world", "favorite_count": 3, "retweet_count": 1,
|
||||
"reply_count": 0, "id_str": "111", "lang": "en"},
|
||||
"core": {"user_results": {"result": {"legacy": {"screen_name": "alice", "name": "Alice"}}}},
|
||||
"views": {"count": "42"},
|
||||
}
|
||||
|
||||
CANNED_TIMELINE = {"data": {"search_by_raw_query": {"search_timeline": {"timeline": {"instructions": [
|
||||
{"type": "TimelineAddEntries", "entries": [
|
||||
{"entryId": "tweet-111", "content": {"itemContent": {"tweet_results": {"result": CANNED_TWEET}}}},
|
||||
{"entryId": "cursor-bottom", "content": {"cursorType": "Bottom", "value": "CURSOR123"}},
|
||||
]},
|
||||
]}}}}}
|
||||
|
||||
|
||||
# -- id extraction + normalizers -------------------------------------------
|
||||
|
||||
def test_tweet_id_from_url_and_bare():
|
||||
assert tweet_id_of("https://x.com/alice/status/1850000000000000123") == "1850000000000000123"
|
||||
assert tweet_id_of("1850000000000000123") == "1850000000000000123"
|
||||
assert tweet_id_of("t_nope") == "t_nope"
|
||||
|
||||
|
||||
def test_normalize_tweet_normalizes():
|
||||
t = normalize_tweet(CANNED_TWEET)
|
||||
assert t["id"] == "111" and t["author"] == "alice" and t["name"] == "Alice"
|
||||
assert t["text"] == "hello world" and t["likes"] == 3 and t["views"] == "42"
|
||||
|
||||
|
||||
def test_normalize_tweet_unwraps_visibility_wrapper():
|
||||
wrapped = {"__typename": "TweetWithVisibilityResults", "tweet": CANNED_TWEET}
|
||||
assert normalize_tweet(wrapped)["id"] == "111"
|
||||
|
||||
|
||||
def test_long_text_truncated():
|
||||
big = {"rest_id": "9", "legacy": {"full_text": "x" * 4000, "id_str": "9"}}
|
||||
body = normalize_tweet(big)["text"]
|
||||
assert len(body) < 4000 and "+2800 chars" in body
|
||||
|
||||
|
||||
# -- rate limiter ----------------------------------------------------------
|
||||
|
||||
def test_first_read_is_prompt():
|
||||
start = time.time()
|
||||
rate_limit.acquire("read")
|
||||
assert time.time() - start < 1.2
|
||||
|
||||
|
||||
def test_429_backoff_delays_next_request():
|
||||
rate_limit.note_response(429, {"retry-after": "1"})
|
||||
start = time.time()
|
||||
rate_limit.acquire("read")
|
||||
assert time.time() - start >= 0.8
|
||||
|
||||
|
||||
# -- dispatch + normalizers (network mocked) -------------------------------
|
||||
|
||||
def test_search_walks_nested_timeline():
|
||||
with patch.object(x_reads, "graphql", return_value=CANNED_TIMELINE):
|
||||
out = handle_tool_call("x_search", {"query": "openswarm", "count": 10})
|
||||
data = json.loads(p_text(out))
|
||||
assert "isError" not in out
|
||||
assert data["cursor"] == "CURSOR123"
|
||||
assert data["tweets"][0]["id"] == "111"
|
||||
assert data["tweets"][0]["author"] == "alice"
|
||||
|
||||
|
||||
def test_like_maps_to_favorite_op():
|
||||
captured: dict = {}
|
||||
|
||||
def fake_graphql(op, variables, **kw):
|
||||
captured["op"], captured["vars"] = op, variables
|
||||
return {}
|
||||
|
||||
with patch.object(x_writes, "graphql", fake_graphql):
|
||||
out = handle_tool_call("x_like", {"target": "https://x.com/a/status/1850000000000000111"})
|
||||
assert captured["op"] == "FavoriteTweet"
|
||||
assert captured["vars"]["tweet_id"] == "1850000000000000111"
|
||||
assert json.loads(p_text(out))["liked"] is True
|
||||
|
||||
|
||||
def test_unlike_maps_to_unfavorite_op():
|
||||
captured: dict = {}
|
||||
with patch.object(x_writes, "graphql", lambda op, v, **k: captured.setdefault("op", op) or {}):
|
||||
handle_tool_call("x_like", {"target": "111", "unlike": True})
|
||||
assert captured["op"] == "UnfavoriteTweet"
|
||||
|
||||
|
||||
def test_follow_resolves_id_and_hits_v11():
|
||||
captured: dict = {}
|
||||
|
||||
def fake_rest(method, path, **kw):
|
||||
captured["path"], captured["form"] = path, kw.get("form")
|
||||
return {}
|
||||
|
||||
with patch.object(x_writes, "resolve_user_id", return_value="999"), \
|
||||
patch.object(x_writes, "rest", fake_rest):
|
||||
out = handle_tool_call("x_follow", {"username": "@bob"})
|
||||
assert captured["path"] == "1.1/friendships/create.json"
|
||||
assert captured["form"]["user_id"] == "999"
|
||||
assert json.loads(p_text(out))["following"] is True
|
||||
|
||||
|
||||
def test_session_unavailable_is_actionable():
|
||||
def boom(*a, **k):
|
||||
raise SessionUnavailable("Not logged in to x.com. Open x.com in the OpenSwarm browser, sign in, then retry.")
|
||||
|
||||
with patch.object(x_reads, "rest", boom):
|
||||
out = handle_tool_call("x_whoami", {})
|
||||
assert out.get("isError") is True
|
||||
assert "logged in" in p_text(out).lower()
|
||||
|
||||
|
||||
def test_unknown_tool_errors():
|
||||
out = handle_tool_call("x_nonsense", {})
|
||||
assert out.get("isError") is True
|
||||
Reference in New Issue
Block a user