#!/usr/bin/env python3
"""Dependency-free Crosspost CLI and stdio MCP bridge for finished local videos.

Credentials stay in CROSSPOST_API_KEY. Private durable jobs protect against
duplicate upload initialization and post creation across process restarts.
"""

import argparse
import datetime as dt
import fcntl
import hashlib
import http.client
import ipaddress
import json
import os
from pathlib import Path
import re
import socket
import ssl
import stat
import sys
import tempfile
import time
from urllib.parse import urlsplit
import uuid

VERSION = "1.0.0"
API_BASE = "https://cross-post.app/api/v1"
MIME_TYPES = {".mp4": "video/mp4", ".mov": "video/quicktime", ".webm": "video/webm"}
MAX_RESPONSE = 2 * 1024 * 1024
JOB_PATTERN = re.compile(r"[A-Za-z0-9][A-Za-z0-9_-]{0,79}\Z")
UTC = dt.timezone.utc
PUBLIC_ERROR_CODES = frozenset({
    "authentication_required", "invalid_api_key", "invalid_scope", "insufficient_scope", "forbidden",
    "subscription_required", "api_access_disabled", "rate_limit_exceeded", "rate_limited", "validation_error",
    "idempotency_required", "idempotency_in_progress", "idempotency_mismatch", "idempotency_outcome_unknown",
    "account_deletion_in_progress", "internal_error", "reconciliation_required", "social_publish_error",
    "unsupported_fanout", "limit_exceeded", "not_found", "method_not_allowed", "media_processing",
})


class AgentError(Exception):
    def __init__(self, code, message, status=None, job_id=None):
        super().__init__(message)
        self.code, self.message, self.status, self.job_id = code, message, status, job_id

    def public(self):
        result = {"code": self.code, "message": self.message}
        if self.status is not None:
            result["status"] = self.status
        if self.job_id is not None:
            result["job_id"] = self.job_id
        return {"error": result}


def encoded_json(value):
    return json.dumps(value, ensure_ascii=False, separators=(",", ":"), sort_keys=True).encode("utf-8")


def positive_integer(value, label):
    if type(value) is not int or value <= 0 or value > 2147483647:
        raise AgentError("invalid_arguments", label + " must be a positive integer.")
    return value


def checked_url(value, *, base=False):
    if not isinstance(value, str) or len(value) > 16384 or "\\" in value or re.search(r"[\x00-\x20\x7f]", value):
        raise AgentError("unsafe_url", "The API or upload URL is invalid.")
    try:
        parsed = urlsplit(value)
        port = parsed.port
    except ValueError:
        raise AgentError("unsafe_url", "The API or upload URL is invalid.") from None
    host = parsed.hostname
    if parsed.scheme != "https" or not host or parsed.username is not None or parsed.password is not None:
        raise AgentError("unsafe_url", "Only credential-free HTTPS URLs are allowed.")
    if parsed.fragment or port not in (None, 443) or (base and parsed.query):
        raise AgentError("unsafe_url", "The API or upload URL is invalid.")
    host = host.rstrip(".").lower()
    if host == "localhost" or host.endswith((".localhost", ".local", ".internal")):
        raise AgentError("unsafe_url", "Private network destinations are not allowed.")
    try:
        address = ipaddress.ip_address(host)
    except ValueError:
        address = None
    if address is not None and not address.is_global:
        raise AgentError("unsafe_url", "Private network destinations are not allowed.")
    return parsed


class PublicHTTPSConnection(http.client.HTTPSConnection):
    """Resolve once, refuse all non-public addresses, and pin the TLS socket."""

    def connect(self):
        addresses = socket.getaddrinfo(self.host, self.port, type=socket.SOCK_STREAM)
        if not addresses or any(not ipaddress.ip_address(item[4][0]).is_global for item in addresses):
            raise AgentError("unsafe_url", "Private network destinations are not allowed.")
        last_error = None
        for family, socktype, protocol, _, address in addresses:
            raw = socket.socket(family, socktype, protocol)
            raw.settimeout(self.timeout)
            try:
                raw.connect(address)
                self.sock = self._context.wrap_socket(raw, server_hostname=self.host)
                return
            except (OSError, ssl.SSLError) as error:
                raw.close()
                last_error = error
        raise last_error or OSError("connection failed")


class HTTPTransport:
    def request(self, method, url, headers, body=None, *, size=None, expected_sha256=None):
        parsed = checked_url(url)
        connection = PublicHTTPSConnection(parsed.hostname, timeout=300 if method == "PUT" else 30,
                                           context=ssl.create_default_context())
        path = parsed.path or "/"
        if parsed.query:
            path += "?" + parsed.query
        try:
            if hasattr(body, "read"):
                # http.client sends an iterable without retaining the video in memory.
                digest = hashlib.sha256()
                transferred = [0]

                def chunks():
                    while True:
                        chunk = body.read(1024 * 1024)
                        if not chunk:
                            break
                        transferred[0] += len(chunk)
                        if transferred[0] > size:
                            raise AgentError("file_changed", "The video changed during upload; no post was created.")
                        digest.update(chunk)
                        yield chunk

                headers = dict(headers, **{"Content-Length": str(size)})
                connection.request(method, path, body=chunks(), headers=headers)
                if transferred[0] != size or digest.hexdigest() != expected_sha256:
                    raise AgentError("file_changed", "The video changed during upload; no post was created.")
            else:
                connection.request(method, path, body=body, headers=headers)
            response = connection.getresponse()
            payload = response.read(MAX_RESPONSE + 1)
            if len(payload) > MAX_RESPONSE:
                raise AgentError("invalid_response", "The server response exceeded the size limit.")
            # Redirects are returned as errors, never followed with credentials/file bytes.
            return response.status, payload, {key.lower(): value for key, value in response.getheaders()}
        except AgentError:
            raise
        except (OSError, ssl.SSLError, http.client.HTTPException):
            raise AgentError("transport_error", "Network outcome is uncertain. Resume the same job; do not create a replacement job.") from None
        finally:
            connection.close()


def safe_http_error(status, payload):
    error = payload.get("error", {}) if isinstance(payload, dict) else {}
    code = error.get("code") if isinstance(error, dict) else None
    if not isinstance(code, str) or code not in PUBLIC_ERROR_CODES:
        code = "http_error"
    return AgentError(code, "Crosspost returned HTTP " + str(status) + ". Check the account, key scopes, and saved job before retrying.", status)


class JobStore:
    def __init__(self, directory):
        self.directory = Path(directory).expanduser().absolute()

    def lock(self, job_id):
        self.directory.mkdir(mode=0o700, parents=True, exist_ok=True)
        info = self.directory.lstat()
        if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid() or stat.S_IMODE(info.st_mode) & 0o077:
            raise AgentError("unsafe_state_directory", "CROSSPOST_STATE_DIR must be an owned, private directory (mode 0700).")
        lock_path = self.directory / (job_id + ".lock")
        descriptor = os.open(lock_path, os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_NONBLOCK, 0o600)
        info = os.fstat(descriptor)
        if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() or stat.S_IMODE(info.st_mode) & 0o077:
            os.close(descriptor)
            raise AgentError("unsafe_state_file", "A job state file has unsafe ownership or permissions.")
        try:
            fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB)
        except BlockingIOError:
            os.close(descriptor)
            raise AgentError("job_busy", "This job is already running. Wait before resuming the same job.") from None
        return os.fdopen(descriptor, "r+")

    def read(self, job_id):
        path = self.directory / (job_id + ".json")
        try:
            descriptor = os.open(path, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK)
        except FileNotFoundError:
            return None
        with os.fdopen(descriptor, "rb") as stream:
            info = os.fstat(stream.fileno())
            if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() or stat.S_IMODE(info.st_mode) & 0o077:
                raise AgentError("unsafe_state_file", "A job state file has unsafe ownership or permissions.")
            try:
                data = stream.read(MAX_RESPONSE + 1)
                if len(data) > MAX_RESPONSE:
                    raise ValueError()
                value = json.loads(data)
            except (ValueError, UnicodeDecodeError):
                raise AgentError("invalid_state", "The saved job is unreadable. Preserve it for manual recovery.") from None
            if not isinstance(value, dict) or value.get("version") != 1:
                raise AgentError("invalid_state", "The saved job has an unsupported format. Preserve it for manual recovery.")
            return value

    def save(self, job_id, value):
        descriptor, name = tempfile.mkstemp(prefix=".job-", dir=self.directory)
        try:
            with os.fdopen(descriptor, "wb") as stream:
                os.fchmod(stream.fileno(), 0o600)
                stream.write(encoded_json(value))
                stream.flush()
                os.fsync(stream.fileno())
            os.replace(name, self.directory / (job_id + ".json"))
            descriptor = os.open(self.directory, os.O_RDONLY)
            try:
                os.fsync(descriptor)
            finally:
                os.close(descriptor)
        finally:
            if os.path.exists(name):
                os.unlink(name)


def normalized_timestamp(value):
    if not isinstance(value, str) or not re.fullmatch(r"\d{4}-\d{2}-\d{2}T\d{2}:\d{2}(?::\d{2}(?:\.\d{1,6})?)?(?:Z|[+-]\d{2}:\d{2})", value):
        raise AgentError("invalid_arguments", "scheduled_at must be an ISO 8601 datetime with explicit Z or UTC offset.")
    try:
        timestamp = dt.datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(UTC)
    except ValueError:
        raise AgentError("invalid_arguments", "scheduled_at is not a valid datetime.") from None
    return timestamp, timestamp.isoformat(timespec="microseconds" if timestamp.microsecond else "seconds").replace("+00:00", "Z")


def valid_token(value):
    return isinstance(value, str) and 0 < len(value) <= 2048 and not re.search(r"[\x00-\x20\x7f]", value)


def public_post(value):
    if not isinstance(value, dict):
        raise AgentError("invalid_response", "Crosspost did not return a valid post record.")
    positive_integer(value.get("id"), "Post ID")
    fields = ("id", "status", "scheduled_at", "published_at", "created_at")
    result = {name: value[name] for name in fields if name in value}
    return result


class CrosspostAgent:
    def __init__(self, api_key=None, base_url=None, state_dir=None, transport=None, sleep=time.sleep, now=None):
        self.api_key = api_key if api_key is not None else os.environ.get("CROSSPOST_API_KEY", "")
        if not isinstance(self.api_key, str) or not self.api_key or len(self.api_key) > 512 or re.search(r"[\x00-\x20\x7f]", self.api_key):
            raise AgentError("missing_api_key", "Set CROSSPOST_API_KEY to your Crosspost Developer API key.")
        self.base_url = (base_url or os.environ.get("CROSSPOST_API_BASE_URL", API_BASE)).rstrip("/")
        checked_url(self.base_url, base=True)
        default_state = Path.home() / ".local" / "state" / "crosspost-agent"
        self.store = JobStore(state_dir or os.environ.get("CROSSPOST_STATE_DIR", str(default_state)))
        self.transport = transport or HTTPTransport()
        self.sleep = sleep
        self.now = now or (lambda: dt.datetime.now(UTC))
        self.key_fingerprint = hashlib.sha256(self.api_key.encode()).hexdigest()

    def request(self, method, path, body=None, *, raw=None, idempotency_key=None, safe_retry=False):
        headers = {"Authorization": "Bearer " + self.api_key, "Accept": "application/json", "User-Agent": "crosspost-agent/" + VERSION}
        if body is not None or raw is not None:
            headers["Content-Type"] = "application/json"
        if idempotency_key:
            headers["X-Idempotency-Key"] = idempotency_key
        payload = raw if raw is not None else (encoded_json(body) if body is not None else None)
        for attempt in range(2 if safe_retry else 1):
            status, response, response_headers = self.transport.request(method, self.base_url + path, headers, payload)
            try:
                result = json.loads(response)
            except (ValueError, UnicodeDecodeError):
                raise AgentError("invalid_response", "Crosspost returned an unreadable response. Preserve the same job for recovery.") from None
            if not isinstance(result, dict):
                raise AgentError("invalid_response", "Crosspost returned an invalid response.")
            if safe_retry and attempt == 0 and status in (429, 503):
                delay = response_headers.get("retry-after", "2")
                self.sleep(min(10, max(1, int(delay))) if str(delay).isdigit() else 2)
                continue
            return status, result

    def get_data(self, path):
        status, result = self.request("GET", path, safe_retry=True)
        if not 200 <= status < 300:
            raise safe_http_error(status, result)
        if "data" not in result:
            raise AgentError("invalid_response", "Crosspost response is missing data.")
        return result["data"]

    def accounts(self):
        value = self.get_data("/accounts")
        if not isinstance(value, list) or any(not isinstance(row, dict) for row in value):
            raise AgentError("invalid_response", "Crosspost did not return an account list.")
        allowed = ("id", "platform", "username", "display_name", "is_active", "provider", "account_group_id", "authorization_status", "authorization_message")
        return [{name: row[name] for name in allowed if name in row} for row in value]

    def usage(self):
        value = self.get_data("/usage")
        if not isinstance(value, dict):
            raise AgentError("invalid_response", "Crosspost did not return usage information.")
        allowed = ("subscription_tier", "subscription_expires_at", "provider_write_allowed", "posts", "scheduled_posts", "social_accounts", "max_media_size_mb")
        return {name: value[name] for name in allowed if name in value}

    def get_post(self, post_id):
        return public_post(self.get_data("/posts/" + str(positive_integer(post_id, "post_id"))))

    def preflight(self, intent):
        usage = self.usage()
        if usage.get("provider_write_allowed") is not True:
            raise AgentError("subscription_required", "An active publishing entitlement is required before uploading or scheduling.")
        maximum = usage.get("max_media_size_mb")
        if type(maximum) not in (int, float) or maximum <= 0:
            raise AgentError("invalid_response", "Crosspost did not return a valid media size limit.")
        if intent["size_bytes"] > maximum * 1024 * 1024:
            raise AgentError("file_too_large", "The video exceeds this plan's media size limit.")
        # The schedule API reserves scheduled_posts_this_month, not the separate
        # immediate-publication counter. Do not reject a valid scheduling plan
        # just because immediate-post allowance is exhausted.
        allowance = usage.get("scheduled_posts", {})
        remaining = allowance.get("remaining") if isinstance(allowance, dict) else None
        if type(remaining) is not int or remaining == 0 or remaining < -1:
            raise AgentError("schedule_limit", "Scheduled-post allowance is unavailable or exhausted.")
        capabilities = self.get_data("/capabilities")
        if not isinstance(capabilities, dict):
            raise AgentError("invalid_response", "Crosspost capability metadata is unavailable.")
        posting, media = capabilities.get("posting", {}), capabilities.get("media", {})
        if not isinstance(posting, dict) or type(posting.get("same_account_group_required")) is not bool:
            raise AgentError("invalid_response", "Crosspost account-group capability metadata is unavailable.")
        if capabilities.get("provider_write_allowed") is not True or posting.get("enabled") is not True or posting.get("scheduling") is not True:
            raise AgentError("scheduling_unavailable", "This provider or entitlement does not currently allow scheduling.")
        if not isinstance(media, dict) or media.get("upload_method") != "PUT" or media.get("confirmation_required") is not True:
            raise AgentError("unsupported_upload_flow", "This connector requires confirmed HTTPS presigned PUT uploads, as provided by production Bundle.")
        allowed_types = media.get("allowed_content_types") if isinstance(media, dict) else None
        if not isinstance(allowed_types, list) or not all(isinstance(item, str) for item in allowed_types):
            raise AgentError("invalid_response", "Crosspost media capability metadata is unavailable.")
        if intent["content_type"] not in allowed_types:
            raise AgentError("unsupported_media", "This provider does not support that video format. Export MP4 or MOV instead.")
        rows = self.accounts()
        selected = []
        for account_id in intent["account_ids"]:
            matches = [row for row in rows if row.get("id") == account_id and row.get("is_active") is True]
            if len(matches) != 1:
                raise AgentError("invalid_account", "A selected account is missing or inactive.")
            account = matches[0]
            if account.get("authorization_status") == "reconnect_required":
                raise AgentError("reconnect_required", "Reauthorize the selected social account in Crosspost before scheduling.")
            if account.get("authorization_status") != "unknown":
                raise AgentError("invalid_response", "Account authorization metadata is unavailable.")
            if account.get("platform") == "pinterest":
                raise AgentError("unsupported_platform", "Pinterest needs a board selection. Use the Crosspost dashboard for Pinterest; this video connector does not yet supply board options.")
            selected.append(account)
        if posting["same_account_group_required"]:
            group_ids = [row.get("account_group_id") for row in selected]
            if any(type(group_id) is not int or group_id <= 0 for group_id in group_ids) or len(set(group_ids)) != 1:
                raise AgentError("mixed_account_groups", "Choose accounts from one account group. Each Bundle video post must stay within one group.")

    def schedule_video(self, file_path, account_ids, scheduled_at, caption="", job_id=None):
        if not isinstance(job_id, str) or not JOB_PATTERN.fullmatch(job_id):
            raise AgentError("invalid_arguments", "A stable job_id is required: use 1–80 letters, digits, underscores or hyphens and reuse it for retries.")
        try:
            return self._schedule_video(file_path, account_ids, scheduled_at, caption, job_id)
        except AgentError as error:
            error.job_id = job_id
            raise
        except (OSError, ValueError, KeyError, TypeError):
            raise AgentError("local_error", "The video or private job state could not be processed. Preserve the job for recovery.", job_id=job_id) from None

    def _schedule_video(self, file_path, account_ids, scheduled_at, caption, job_id):
        if not isinstance(file_path, str) or not file_path or "\x00" in file_path:
            raise AgentError("invalid_arguments", "file_path must name a finished local video.")
        if not isinstance(account_ids, list) or not account_ids or len(account_ids) > 50:
            raise AgentError("invalid_arguments", "account_ids must be a non-empty list of positive integers.")
        ids = [positive_integer(value, "account_ids") for value in account_ids]
        if len(set(ids)) != len(ids):
            raise AgentError("invalid_arguments", "account_ids must not contain duplicates.")
        if not isinstance(caption, str) or len(caption.encode("utf-8")) > 10000:
            raise AgentError("invalid_arguments", "caption must be a string of at most 10,000 UTF-8 bytes.")
        timestamp, canonical_time = normalized_timestamp(scheduled_at)
        path = Path(file_path).expanduser().resolve(strict=True)
        content_type = MIME_TYPES.get(path.suffix.lower())
        if content_type is None:
            raise AgentError("unsupported_media", "Use a finished local MP4, MOV or supported WebM video.")
        if len(path.name.encode("utf-8")) > 200 or re.search(r"[\x00-\x1f\x7f]", path.name):
            raise AgentError("invalid_arguments", "The video filename is too long or invalid.")
        descriptor = os.open(path, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK)
        with os.fdopen(descriptor, "rb") as video, self.store.lock(job_id):
            info = os.fstat(video.fileno())
            if not stat.S_ISREG(info.st_mode) or info.st_size <= 0:
                raise AgentError("invalid_arguments", "The video must be a non-empty regular file.")
            digest = hashlib.sha256()
            hashed_bytes = 0
            for chunk in iter(lambda: video.read(1024 * 1024), b""):
                digest.update(chunk)
                hashed_bytes += len(chunk)
            after_hash = os.fstat(video.fileno())
            if hashed_bytes != info.st_size or (after_hash.st_size, after_hash.st_mtime_ns, after_hash.st_ctime_ns) != (info.st_size, info.st_mtime_ns, info.st_ctime_ns):
                raise AgentError("file_changed", "The video changed while preparing the job. Finish rendering before scheduling.")
            video.seek(0)
            intent = {"file_path": str(path), "filename": path.name, "size_bytes": info.st_size,
                      "sha256": digest.hexdigest(), "content_type": content_type, "caption": caption,
                      "account_ids": sorted(ids), "scheduled_at": canonical_time,
                      "api_base_url": self.base_url, "api_key_fingerprint": self.key_fingerprint}
            job = self.store.read(job_id)
            if job is not None:
                if job.get("intent") != intent:
                    raise AgentError("intent_changed", "This job belongs to a different video, caption, accounts, time, API key or API origin. Do not replace its saved intent.")
                if job.get("state") == "completed":
                    return job["receipt"]
                if job.get("state") == "failed":
                    raise AgentError("completed_failure", "This job has a terminal failure. Preserve it; do not silently create a replacement.")
                if job.get("state") == "init_started":
                    raise AgentError("upload_initialization_unknown", "Upload initialization has an unknown outcome. Manual recovery is required; it will not be repeated.")
                if job.get("state") == "reconciliation_required":
                    return self.reconciliation_receipt(job_id, job)
            else:
                if timestamp <= self.now():
                    raise AgentError("past_schedule", "Choose a future scheduled_at datetime.")
                self.preflight(intent)
                job = {"version": 1, "job_id": job_id, "intent": intent, "state": "init_started",
                       "idempotency_key": "cp-agent-" + uuid.uuid4().hex}
                self.store.save(job_id, job)
                status, result = self.request("POST", "/media/upload", {"filename": path.name, "content_type": content_type,
                                                                      "size_bytes": info.st_size, "destination_account_id": sorted(ids)[0]})
                if not 200 <= status < 300:
                    # Even an HTTP failure is retained: there is no init idempotency contract.
                    raise safe_http_error(status, result)
                upload = result.get("data", {})
                if not isinstance(upload, dict) or not valid_token(upload.get("media_id")) or type(upload.get("requires_confirm")) is not bool:
                    raise AgentError("invalid_response", "Upload initialization returned incomplete metadata. Manual recovery is required.")
                if upload.get("upload_method") != "PUT" or upload.get("upload_headers") != {"Content-Type": content_type} or upload["requires_confirm"] is not True:
                    raise AgentError("invalid_response", "Upload initialization returned an unsupported method, headers or confirmation flow. Manual recovery is required.")
                checked_url(upload.get("upload_url"))
                job.update(state="initialized", upload_url=upload["upload_url"], media_id=upload["media_id"],
                           public_url=upload.get("public_url"), requires_confirm=upload["requires_confirm"])
                self.store.save(job_id, job)
            state = job.get("state")
            if state not in ("initialized", "put_started", "uploaded", "confirming", "ready", "create_started"):
                raise AgentError("invalid_state", "The saved job state is unsupported. Preserve it for recovery.")
            if state != "create_started":
                if timestamp <= self.now():
                    raise AgentError("past_schedule", "The scheduled time has passed. Preserve this job for manual recovery; its intent cannot be changed.")
                self.preflight(intent)
            if job["state"] in ("initialized", "put_started"):
                job["state"] = "put_started"
                self.store.save(job_id, job)
                status, _, _ = self.transport.request("PUT", job["upload_url"], {"Content-Type": content_type},
                                                       video, size=info.st_size, expected_sha256=intent["sha256"])
                if not 200 <= status < 300:
                    raise AgentError("upload_failed", "The presigned upload failed. Resume only this same job and immutable video; do not initialize a replacement upload.", status)
                job["state"] = "uploaded"
                self.store.save(job_id, job)
            if job["state"] == "uploaded":
                job["state"] = "confirming" if job["requires_confirm"] else "ready"
                self.store.save(job_id, job)
            if job["state"] == "confirming":
                for attempt in range(10):
                    status, result = self.request("POST", "/media/confirm", {"media_id": job["media_id"], "size_bytes": info.st_size}, safe_retry=True)
                    if status not in (200, 202):
                        raise safe_http_error(status, result)
                    confirmed = result.get("data", {})
                    if not isinstance(confirmed, dict) or not valid_token(confirmed.get("media_id")):
                        raise AgentError("invalid_response", "Media confirmation returned an invalid token. Preserve this job for recovery.")
                    # Persist token rotation before any subsequent validation or retry.
                    job["media_id"] = confirmed["media_id"]
                    self.store.save(job_id, job)
                    if type(confirmed.get("ready")) is not bool or type(confirmed.get("requires_confirm")) is not bool:
                        raise AgentError("invalid_response", "Media confirmation returned incomplete readiness metadata.")
                    if status == 200 and confirmed["ready"] and not confirmed["requires_confirm"]:
                        job.update(state="ready", public_url=confirmed.get("public_url"), requires_confirm=False)
                        self.store.save(job_id, job)
                        break
                    if attempt != 9:
                        self.sleep(2)
                else:
                    raise AgentError("media_processing", "The video is still processing. Resume this same job later; its latest media token is saved.")
            if job["state"] == "ready":
                if timestamp <= self.now():
                    raise AgentError("past_schedule", "The scheduled time has passed. Preserve this job for manual recovery.")
                checked_url(job.get("public_url"))
                if len(job["public_url"]) > 512 or len(job["media_id"]) > 255:
                    raise AgentError("invalid_response", "Confirmed media metadata exceeds the post API limits.")
                post_body = {"caption": caption, "destination_account_ids": sorted(ids), "publish_mode": "schedule",
                             "scheduled_at": canonical_time, "timezone": "UTC",
                             "media_urls": [{"url": job["public_url"], "type": content_type, "media_id": job["media_id"]}]}
                serialized_body = encoded_json(post_body)
                job.update(state="create_started", post_body=serialized_body.decode("utf-8"),
                           post_body_sha256=hashlib.sha256(serialized_body).hexdigest())
                self.store.save(job_id, job)
            # Exactly one request per invocation. A resume reuses this same key and byte string.
            serialized_body = job["post_body"].encode("utf-8")
            if hashlib.sha256(serialized_body).hexdigest() != job.get("post_body_sha256") or not re.fullmatch(r"cp-agent-[a-f0-9]{32}", job.get("idempotency_key", "")):
                raise AgentError("invalid_state", "The saved post body or idempotency key is damaged. Preserve this job for recovery.")
            status, result = self.request("POST", "/posts", raw=serialized_body, idempotency_key=job["idempotency_key"])
            if not 200 <= status < 300:
                error_value = result.get("error", {})
                error_post_id = error_value.get("post_id") if isinstance(error_value, dict) else None
                if type(error_post_id) is int and 0 < error_post_id <= 2147483647:
                    job.update(state="reconciliation_required", post_id=error_post_id)
                    self.store.save(job_id, job)
                    return self.reconciliation_receipt(job_id, job)
                if result.get("request_completed") is True:
                    job["state"] = "failed"
                    self.store.save(job_id, job)
                raise safe_http_error(status, result)
            post = public_post(result.get("data"))
            receipt = {"job_id": job_id, "status": "completed", "post": post}
            job.update(state="completed", receipt=receipt)
            # No presigned URLs or media tokens are needed after the completed receipt.
            for name in ("upload_url", "media_id", "public_url", "post_body", "requires_confirm"):
                job.pop(name, None)
            self.store.save(job_id, job)
            return receipt

    def reconciliation_receipt(self, job_id, job):
        # A known local post survives an uncertain provider outcome. Never repeat
        # POST for this job: only read that exact existing post's local status.
        post = self.get_post(positive_integer(job.get("post_id"), "Saved post ID"))
        return {"job_id": job_id, "status": "reconciliation_required", "post": post,
                "message": "An existing post needs provider reconciliation. Do not create a replacement job; this connector only reads its status."}


READ_ANNOTATIONS = {"readOnlyHint": True, "destructiveHint": False, "idempotentHint": True, "openWorldHint": True}
TOOLS = [
    {"name": "list_accounts", "description": "List saved Crosspost social accounts, account groups and known reauthorization warnings.",
     "inputSchema": {"type": "object", "properties": {}, "additionalProperties": False}, "annotations": READ_ANNOTATIONS},
    {"name": "get_usage", "description": "Read publishing entitlement, remaining schedule allowance and media size limit.",
     "inputSchema": {"type": "object", "properties": {}, "additionalProperties": False}, "annotations": READ_ANNOTATIONS},
    {"name": "get_post", "description": "Read a Crosspost post's current local scheduling/publication status; does not force provider refresh.",
     "inputSchema": {"type": "object", "properties": {"post_id": {"type": "integer", "minimum": 1}}, "required": ["post_id"], "additionalProperties": False},
     "annotations": READ_ANNOTATIONS},
    {"name": "schedule_video", "description": "Upload an already-rendered local video and schedule it to explicitly chosen accounts. This writes a real scheduled post: obtain user approval for the exact video, caption, accounts and time. Supply a stable job_id and resume that same job after errors; never create replacement jobs to bypass uncertainty. Video generation is external.",
     "inputSchema": {"type": "object", "properties": {
         "file_path": {"type": "string", "minLength": 1},
         "account_ids": {"type": "array", "items": {"type": "integer", "minimum": 1}, "minItems": 1, "maxItems": 50, "uniqueItems": True},
         "scheduled_at": {"type": "string", "description": "Future ISO 8601 datetime with Z or explicit UTC offset."},
         "caption": {"type": "string", "default": ""},
         "job_id": {"type": "string", "pattern": "^[A-Za-z0-9][A-Za-z0-9_-]{0,79}$"}},
         "required": ["file_path", "account_ids", "scheduled_at", "job_id"], "additionalProperties": False},
     "annotations": {"readOnlyHint": False, "destructiveHint": False, "idempotentHint": False, "openWorldHint": True}},
]


def call_tool(agent, name, arguments):
    definition = next((item for item in TOOLS if item["name"] == name), None)
    if definition is None:
        raise AgentError("unknown_tool", "Unknown Crosspost tool.")
    schema = definition["inputSchema"]
    if not isinstance(arguments, dict) or set(arguments) - set(schema["properties"]) or any(key not in arguments for key in schema.get("required", [])):
        raise AgentError("invalid_arguments", "Tool arguments are missing required fields or contain unsupported fields.")
    if name == "list_accounts":
        return agent.accounts()
    if name == "get_usage":
        return agent.usage()
    if name == "get_post":
        return agent.get_post(**arguments)
    return agent.schedule_video(**arguments)


def serve_mcp(agent_factory=CrosspostAgent, input_stream=sys.stdin, output_stream=sys.stdout):
    initialized = False
    ready = False
    supported_versions = ("2025-11-25", "2025-06-18", "2025-03-26", "2024-11-05")
    agent = None
    for line in input_stream:
        response = None
        request_id = None
        try:
            if len(line.encode("utf-8")) > MAX_RESPONSE:
                raise ValueError()
            request = json.loads(line)
            if not isinstance(request, dict) or request.get("jsonrpc") != "2.0" or not isinstance(request.get("method"), str):
                response = {"error": {"code": -32600, "message": "Invalid request"}}
            else:
                request_id = request.get("id")
                if "id" in request and (type(request_id) not in (str, int) or isinstance(request_id, bool)):
                    response = {"error": {"code": -32600, "message": "Invalid request ID"}}
                    request_id = None
                elif request["method"] == "notifications/initialized" and "id" not in request:
                    ready = initialized
                elif "id" not in request:
                    # Cancellation/unknown notifications never invoke a write tool.
                    continue
                elif request["method"] == "ping":
                    response = {"result": {}}
                elif request["method"] == "initialize":
                    params = request.get("params", {})
                    if initialized or not isinstance(params, dict) or not isinstance(params.get("protocolVersion"), str) or not isinstance(params.get("capabilities"), dict) or not isinstance(params.get("clientInfo"), dict):
                        response = {"error": {"code": -32602, "message": "Invalid initialization"}}
                    else:
                        initialized = True
                        version = params["protocolVersion"] if params["protocolVersion"] in supported_versions else supported_versions[0]
                        response = {"result": {"protocolVersion": version, "capabilities": {"tools": {}},
                                               "serverInfo": {"name": "crosspost-agent", "version": VERSION},
                                               "instructions": "Generate video with your own tools, then use explicit authorized scheduling. Always reuse the same job_id after uncertain outcomes. Never publish without the user's approval."}}
                elif not ready:
                    response = {"error": {"code": -32002, "message": "Initialize and send notifications/initialized first"}}
                elif request["method"] == "tools/list":
                    response = {"result": {"tools": TOOLS}}
                elif request["method"] == "tools/call":
                    params = request.get("params", {})
                    try:
                        if not isinstance(params, dict) or not isinstance(params.get("name"), str):
                            raise AgentError("invalid_arguments", "Tool name and object arguments are required.")
                        if agent is None:
                            agent = agent_factory()
                        result = call_tool(agent, params["name"], params.get("arguments", {}))
                        response = {"result": {"content": [{"type": "text", "text": encoded_json(result).decode("utf-8")}], "isError": False}}
                    except AgentError as error:
                        response = {"result": {"content": [{"type": "text", "text": encoded_json(error.public()).decode("utf-8")}], "isError": True}}
                else:
                    response = {"error": {"code": -32601, "message": "Method not found"}}
        except (ValueError, UnicodeDecodeError):
            response = {"error": {"code": -32700, "message": "Parse error"}}
        except (OSError, TypeError, KeyError):
            response = {"error": {"code": -32603, "message": "Internal error; preserve any active job"}}
        if response is not None:
            response.update(jsonrpc="2.0", id=request_id)
            output_stream.write(encoded_json(response).decode("utf-8") + "\n")
            output_stream.flush()


def main(argv=None):
    parser = argparse.ArgumentParser(description="Crosspost finished-video scheduler and stdio MCP connector")
    commands = parser.add_subparsers(dest="command", required=True)
    commands.add_parser("accounts", help="List connected accounts and groups")
    commands.add_parser("usage", help="Read entitlement and remaining allowance")
    commands.add_parser("mcp", help="Serve MCP tools over stdin/stdout")
    status = commands.add_parser("status", help="Read post status")
    status.add_argument("post_id", type=int)
    schedule = commands.add_parser("schedule-video", help="Upload a finished local video and schedule it")
    schedule.add_argument("--file", required=True, dest="file_path")
    schedule.add_argument("--accounts", required=True, type=int, nargs="+", dest="account_ids")
    schedule.add_argument("--at", required=True, dest="scheduled_at", help="Future ISO 8601 datetime with UTC offset")
    schedule.add_argument("--caption", default="")
    schedule.add_argument("--job-id", required=True, dest="job_id", help="Stable unique ID: reuse it for recovery")
    arguments = parser.parse_args(argv)
    if arguments.command == "mcp":
        serve_mcp()
        return 0
    try:
        agent = CrosspostAgent()
        if arguments.command == "accounts":
            result = agent.accounts()
        elif arguments.command == "usage":
            result = agent.usage()
        elif arguments.command == "status":
            result = agent.get_post(arguments.post_id)
        else:
            result = agent.schedule_video(arguments.file_path, arguments.account_ids, arguments.scheduled_at,
                                          arguments.caption, arguments.job_id)
        print(encoded_json(result).decode("utf-8"))
        return 0
    except AgentError as error:
        print(encoded_json(error.public()).decode("utf-8"), file=sys.stderr)
        return 1


if __name__ == "__main__":
    sys.exit(main())
