Compare commits
11
Commits
a2b859dd89
...
604fb38937
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
604fb38937 | ||
|
|
93e3de3e84 | ||
|
|
2b45db5f5e | ||
|
|
23d4b1fc79 | ||
|
|
9824c5d6bf | ||
|
|
329df8afc3 | ||
|
|
c88ea40459 | ||
|
|
68c40055e7 | ||
|
|
83a3c482ec | ||
|
|
af796d8583 | ||
|
|
3d4a1cb791 |
@@ -0,0 +1,21 @@
|
|||||||
|
MIT License
|
||||||
|
|
||||||
|
Copyright (c) 2026 Kacper Kwapisz
|
||||||
|
|
||||||
|
Permission is hereby granted, free of charge, to any person obtaining a copy
|
||||||
|
of this software and associated documentation files (the "Software"), to deal
|
||||||
|
in the Software without restriction, including without limitation the rights
|
||||||
|
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
|
||||||
|
copies of the Software, and to permit persons to whom the Software is
|
||||||
|
furnished to do so, subject to the following conditions:
|
||||||
|
|
||||||
|
The above copyright notice and this permission notice shall be included in all
|
||||||
|
copies or substantial portions of the Software.
|
||||||
|
|
||||||
|
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||||
|
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
|
||||||
|
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
|
||||||
|
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
|
||||||
|
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
|
||||||
|
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
|
||||||
|
SOFTWARE.
|
||||||
+204
-36
@@ -265,6 +265,15 @@ def resolve_account(accounts: list[dict], account_id: str) -> dict:
|
|||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
|
# RFC 2971 IMAP ID — required by some providers (e.g. netease 163.com / 126.com / yeah.net)
|
||||||
|
# which reject clients that don't identify themselves after login.
|
||||||
|
_IMAP_CLIENT_ID = {
|
||||||
|
"name": "poke-mail",
|
||||||
|
"version": "1.0.0",
|
||||||
|
"vendor": "Poke Interactions",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
def get_imap_client(account: dict) -> IMAPClient:
|
def get_imap_client(account: dict) -> IMAPClient:
|
||||||
port = account["imap_port"]
|
port = account["imap_port"]
|
||||||
use_ssl = port == 993
|
use_ssl = port == 993
|
||||||
@@ -272,6 +281,13 @@ def get_imap_client(account: dict) -> IMAPClient:
|
|||||||
if not use_ssl:
|
if not use_ssl:
|
||||||
client.starttls()
|
client.starttls()
|
||||||
client.login(account["imap_username"], account["imap_password"])
|
client.login(account["imap_username"], account["imap_password"])
|
||||||
|
# Send IMAP ID if the server supports it. Non-fatal: not all servers do,
|
||||||
|
# and we never want this to break an otherwise-working login.
|
||||||
|
try:
|
||||||
|
if client.has_capability("ID"):
|
||||||
|
client.id_(_IMAP_CLIENT_ID)
|
||||||
|
except Exception as e:
|
||||||
|
logger.debug("IMAP ID command failed for %s: %s", account.get("id"), e)
|
||||||
return client
|
return client
|
||||||
|
|
||||||
|
|
||||||
@@ -439,6 +455,55 @@ async def forward_to_poke(
|
|||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def _format_uid_list(uids: list[int], limit: int = 10) -> str:
|
||||||
|
if not uids:
|
||||||
|
return "[]"
|
||||||
|
shown = ", ".join(str(uid) for uid in uids[:limit])
|
||||||
|
if len(uids) > limit:
|
||||||
|
shown += f", ... (+{len(uids) - limit} more)"
|
||||||
|
return f"[{shown}]"
|
||||||
|
|
||||||
|
|
||||||
|
async def _forward_uid_batch(
|
||||||
|
client: IMAPClient,
|
||||||
|
account: dict,
|
||||||
|
folder: str,
|
||||||
|
webhook_url: str,
|
||||||
|
api_key: str,
|
||||||
|
uids: list[int],
|
||||||
|
) -> None:
|
||||||
|
logger.info(
|
||||||
|
"[%s/%s] Forwarding %d message(s) for UIDs %s",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
len(uids),
|
||||||
|
_format_uid_list(uids),
|
||||||
|
)
|
||||||
|
raw_messages = await asyncio.to_thread(client.fetch, uids, ["RFC822"])
|
||||||
|
for uid in uids:
|
||||||
|
data = raw_messages.get(uid, {})
|
||||||
|
raw = data.get(b"RFC822", b"")
|
||||||
|
if not raw:
|
||||||
|
logger.debug(
|
||||||
|
"[%s/%s] Skipping UID %s because RFC822 payload was empty",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
uid,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
email_data = parse_email_message(raw)
|
||||||
|
await forward_to_poke(email_data, account, webhook_url, api_key)
|
||||||
|
|
||||||
|
if account.get("mark_as_read", False):
|
||||||
|
await asyncio.to_thread(client.set_flags, uids, [b"\\Seen"])
|
||||||
|
logger.debug(
|
||||||
|
"[%s/%s] Marked UIDs %s as Seen",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
_format_uid_list(uids),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
# IDLE watcher
|
# IDLE watcher
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
@@ -454,16 +519,26 @@ async def watch_folder(
|
|||||||
backoff = 5
|
backoff = 5
|
||||||
max_backoff = 60
|
max_backoff = 60
|
||||||
|
|
||||||
|
# Hoisted across reconnects so a transient disconnect (e.g. iCloud's
|
||||||
|
# unsolicited FETCH FLAGS during IDLE) doesn't re-baseline the cursor and
|
||||||
|
# silently drop mail that arrived during the reconnect window. Reset only
|
||||||
|
# on first run or when the server's UIDVALIDITY changes (which means UIDs
|
||||||
|
# have been reassigned and the previous cursor is meaningless).
|
||||||
|
last_seen_uid: Optional[int] = None
|
||||||
|
last_uidvalidity: Optional[int] = None
|
||||||
|
|
||||||
while not stop_event.is_set():
|
while not stop_event.is_set():
|
||||||
client = None
|
client = None
|
||||||
try:
|
try:
|
||||||
client = await asyncio.to_thread(get_imap_client, account)
|
client = await asyncio.to_thread(get_imap_client, account)
|
||||||
await asyncio.to_thread(
|
select_info = await asyncio.to_thread(
|
||||||
client.select_folder,
|
client.select_folder,
|
||||||
folder,
|
folder,
|
||||||
readonly=not account.get("mark_as_read", False),
|
readonly=not account.get("mark_as_read", False),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
uidvalidity = select_info.get(b"UIDVALIDITY") if select_info else None
|
||||||
|
|
||||||
# Check IDLE support
|
# Check IDLE support
|
||||||
if not client.has_capability("IDLE"):
|
if not client.has_capability("IDLE"):
|
||||||
logger.warning(
|
logger.warning(
|
||||||
@@ -471,18 +546,56 @@ async def watch_folder(
|
|||||||
account["id"],
|
account["id"],
|
||||||
folder,
|
folder,
|
||||||
)
|
)
|
||||||
|
# Mutable cursor state — _poll_folder updates it so reconnects
|
||||||
|
# don't re-baseline.
|
||||||
|
cursor = {
|
||||||
|
"last_seen_uid": last_seen_uid,
|
||||||
|
"last_uidvalidity": last_uidvalidity,
|
||||||
|
}
|
||||||
|
try:
|
||||||
await _poll_folder(
|
await _poll_folder(
|
||||||
client, account, folder, webhook_url, api_key, stop_event
|
client,
|
||||||
|
account,
|
||||||
|
folder,
|
||||||
|
webhook_url,
|
||||||
|
api_key,
|
||||||
|
stop_event,
|
||||||
|
cursor,
|
||||||
)
|
)
|
||||||
|
finally:
|
||||||
|
last_seen_uid = cursor["last_seen_uid"]
|
||||||
|
last_uidvalidity = cursor["last_uidvalidity"]
|
||||||
|
# _poll_folder only returns when stop_event is set; on errors it
|
||||||
|
# raises and we fall through to the outer reconnect path.
|
||||||
return
|
return
|
||||||
|
|
||||||
# Record existing unseen UIDs so we only forward truly new ones
|
if last_seen_uid is None or uidvalidity != last_uidvalidity:
|
||||||
existing_unseen = set(await asyncio.to_thread(client.search, ["UNSEEN"]))
|
if last_uidvalidity is not None and uidvalidity != last_uidvalidity:
|
||||||
logger.info(
|
logger.warning(
|
||||||
"[%s/%s] Watching for new emails via IDLE (%d existing unseen skipped)",
|
"[%s/%s] UIDVALIDITY changed (%s → %s) — resetting cursor",
|
||||||
account["id"],
|
account["id"],
|
||||||
folder,
|
folder,
|
||||||
len(existing_unseen),
|
last_uidvalidity,
|
||||||
|
uidvalidity,
|
||||||
|
)
|
||||||
|
mailbox_uids = await asyncio.to_thread(client.search, ["ALL"])
|
||||||
|
last_seen_uid = mailbox_uids[-1] if mailbox_uids else 0
|
||||||
|
last_uidvalidity = uidvalidity
|
||||||
|
logger.info(
|
||||||
|
"[%s/%s] Watching for new emails via IDLE from UID %d (%d existing message(s), mark_as_read=%s)",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
last_seen_uid,
|
||||||
|
len(mailbox_uids),
|
||||||
|
account.get("mark_as_read", False),
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
logger.info(
|
||||||
|
"[%s/%s] Resumed IDLE watch from UID %d (mark_as_read=%s)",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
last_seen_uid,
|
||||||
|
account.get("mark_as_read", False),
|
||||||
)
|
)
|
||||||
backoff = 5
|
backoff = 5
|
||||||
|
|
||||||
@@ -502,25 +615,36 @@ async def watch_folder(
|
|||||||
"[%s/%s] IDLE responses: %s", account["id"], folder, responses
|
"[%s/%s] IDLE responses: %s", account["id"], folder, responses
|
||||||
)
|
)
|
||||||
|
|
||||||
uids = await asyncio.to_thread(client.search, ["UNSEEN"])
|
mailbox_uids = await asyncio.to_thread(client.search, ["ALL"])
|
||||||
# Only forward emails that arrived after we started watching
|
new_uids = [uid for uid in mailbox_uids if uid > last_seen_uid]
|
||||||
new_uids = [u for u in uids if u not in existing_unseen]
|
|
||||||
if not new_uids:
|
if not new_uids:
|
||||||
continue
|
logger.debug(
|
||||||
|
"[%s/%s] No new UIDs above %d (mailbox latest UID %d)",
|
||||||
raw_messages = await asyncio.to_thread(
|
account["id"],
|
||||||
client.fetch, new_uids, ["RFC822"]
|
folder,
|
||||||
|
last_seen_uid,
|
||||||
|
mailbox_uids[-1] if mailbox_uids else last_seen_uid,
|
||||||
)
|
)
|
||||||
for uid, data in raw_messages.items():
|
|
||||||
raw = data.get(b"RFC822", b"")
|
|
||||||
if not raw:
|
|
||||||
continue
|
continue
|
||||||
email_data = parse_email_message(raw)
|
|
||||||
await forward_to_poke(email_data, account, webhook_url, api_key)
|
|
||||||
|
|
||||||
if account.get("mark_as_read", False):
|
logger.info(
|
||||||
await asyncio.to_thread(client.set_flags, new_uids, [b"\\Seen"])
|
"[%s/%s] Found %d new UID(s) above %d: %s",
|
||||||
existing_unseen.update(new_uids)
|
account["id"],
|
||||||
|
folder,
|
||||||
|
len(new_uids),
|
||||||
|
last_seen_uid,
|
||||||
|
_format_uid_list(new_uids),
|
||||||
|
)
|
||||||
|
await _forward_uid_batch(
|
||||||
|
client, account, folder, webhook_url, api_key, new_uids
|
||||||
|
)
|
||||||
|
last_seen_uid = new_uids[-1]
|
||||||
|
logger.debug(
|
||||||
|
"[%s/%s] Advanced UID cursor to %d",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
last_seen_uid,
|
||||||
|
)
|
||||||
|
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
@@ -549,24 +673,68 @@ async def _poll_folder(
|
|||||||
webhook_url: str,
|
webhook_url: str,
|
||||||
api_key: str,
|
api_key: str,
|
||||||
stop_event: asyncio.Event,
|
stop_event: asyncio.Event,
|
||||||
|
cursor: dict,
|
||||||
):
|
):
|
||||||
"""Fallback polling for servers without IDLE support. Checks every 60 seconds."""
|
"""Fallback polling for servers without IDLE support. Checks every 60 seconds.
|
||||||
logger.info("[%s/%s] Polling for new emails every 60s", account["id"], folder)
|
|
||||||
|
`cursor` is a mutable dict with keys 'last_seen_uid' and 'last_uidvalidity'
|
||||||
|
shared with watch_folder so cursor state survives reconnects.
|
||||||
|
"""
|
||||||
|
if cursor.get("last_seen_uid") is None:
|
||||||
|
mailbox_uids = await asyncio.to_thread(client.search, ["ALL"])
|
||||||
|
cursor["last_seen_uid"] = mailbox_uids[-1] if mailbox_uids else 0
|
||||||
|
logger.info(
|
||||||
|
"[%s/%s] Polling for new emails every 60s from UID %d (%d existing message(s), mark_as_read=%s)",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
cursor["last_seen_uid"],
|
||||||
|
len(mailbox_uids),
|
||||||
|
account.get("mark_as_read", False),
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
logger.info(
|
||||||
|
"[%s/%s] Resumed polling from UID %d (mark_as_read=%s)",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
cursor["last_seen_uid"],
|
||||||
|
account.get("mark_as_read", False),
|
||||||
|
)
|
||||||
while not stop_event.is_set():
|
while not stop_event.is_set():
|
||||||
try:
|
try:
|
||||||
uids = await asyncio.to_thread(client.search, ["UNSEEN"])
|
last_seen_uid = cursor["last_seen_uid"]
|
||||||
if uids:
|
mailbox_uids = await asyncio.to_thread(client.search, ["ALL"])
|
||||||
raw_messages = await asyncio.to_thread(client.fetch, uids, ["RFC822"])
|
new_uids = [uid for uid in mailbox_uids if uid > last_seen_uid]
|
||||||
for uid, data in raw_messages.items():
|
if new_uids:
|
||||||
raw = data.get(b"RFC822", b"")
|
logger.info(
|
||||||
if not raw:
|
"[%s/%s] Found %d new UID(s) above %d: %s",
|
||||||
continue
|
account["id"],
|
||||||
email_data = parse_email_message(raw)
|
folder,
|
||||||
await forward_to_poke(email_data, account, webhook_url, api_key)
|
len(new_uids),
|
||||||
if account.get("mark_as_read", False):
|
last_seen_uid,
|
||||||
await asyncio.to_thread(client.set_flags, uids, [b"\\Seen"])
|
_format_uid_list(new_uids),
|
||||||
|
)
|
||||||
|
await _forward_uid_batch(
|
||||||
|
client, account, folder, webhook_url, api_key, new_uids
|
||||||
|
)
|
||||||
|
cursor["last_seen_uid"] = new_uids[-1]
|
||||||
|
logger.debug(
|
||||||
|
"[%s/%s] Advanced UID cursor to %d",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
cursor["last_seen_uid"],
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
logger.debug(
|
||||||
|
"[%s/%s] No new UIDs above %d (mailbox latest UID %d)",
|
||||||
|
account["id"],
|
||||||
|
folder,
|
||||||
|
last_seen_uid,
|
||||||
|
mailbox_uids[-1] if mailbox_uids else last_seen_uid,
|
||||||
|
)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("[%s/%s] Poll error: %s", account["id"], folder, e)
|
logger.warning(
|
||||||
|
"[%s/%s] Poll error: %s", account["id"], folder, e
|
||||||
|
)
|
||||||
raise # reconnect via outer loop
|
raise # reconnect via outer loop
|
||||||
await asyncio.sleep(60)
|
await asyncio.sleep(60)
|
||||||
|
|
||||||
|
|||||||
@@ -207,10 +207,23 @@ PYEOF
|
|||||||
fi
|
fi
|
||||||
fi
|
fi
|
||||||
|
|
||||||
# 4. MCP_API_KEY — generate once and persist to .env
|
# ── Load .env early so POKE_TUNNEL is visible to the MCP_API_KEY logic below ──
|
||||||
if [ ! -f .env ] \
|
if [ -f .env ]; then
|
||||||
|
set -a
|
||||||
|
source .env
|
||||||
|
set +a
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Tunnel-mode detection must happen BEFORE the MCP_API_KEY block: in tunnel
|
||||||
|
# mode (POKE_TUNNEL=1, the default) the local server runs unauthenticated and
|
||||||
|
# the tunnel handles auth, so we must NOT generate a key — doing so previously
|
||||||
|
# overwrote user-supplied keys and broke working setups (issue #9).
|
||||||
|
POKE_TUNNEL="${POKE_TUNNEL:-1}"
|
||||||
|
|
||||||
|
# 4. MCP_API_KEY — generate once and persist to .env (skipped in tunnel mode)
|
||||||
|
if [ "${POKE_TUNNEL}" != "1" ] && { [ ! -f .env ] \
|
||||||
|| grep -Eq '^[[:space:]]*MCP_API_KEY=your-secret-key-here' .env 2>/dev/null \
|
|| grep -Eq '^[[:space:]]*MCP_API_KEY=your-secret-key-here' .env 2>/dev/null \
|
||||||
|| ! grep -Eq '^[[:space:]]*MCP_API_KEY=.+' .env 2>/dev/null; then
|
|| ! grep -Eq '^[[:space:]]*MCP_API_KEY=.+' .env 2>/dev/null; }; then
|
||||||
RANDOM_KEY=$(python3 -c "
|
RANDOM_KEY=$(python3 -c "
|
||||||
import secrets, string
|
import secrets, string
|
||||||
alphabet = string.ascii_letters + string.digits
|
alphabet = string.ascii_letters + string.digits
|
||||||
@@ -233,18 +246,13 @@ PYEOF
|
|||||||
fi
|
fi
|
||||||
echo " ✓ MCP_API_KEY generated and saved to .env"
|
echo " ✓ MCP_API_KEY generated and saved to .env"
|
||||||
echo ""
|
echo ""
|
||||||
fi
|
# Re-source so the freshly generated key is visible to the rest of the script
|
||||||
|
|
||||||
# ── Load .env ─────────────────────────────────────────────────────────────────
|
|
||||||
if [ -f .env ]; then
|
|
||||||
set -a
|
set -a
|
||||||
source .env
|
source .env
|
||||||
set +a
|
set +a
|
||||||
fi
|
fi
|
||||||
|
|
||||||
# ── Tunnel-mode detection ─────────────────────────────────────────────────────
|
# ── Tunnel-mode enforcement ───────────────────────────────────────────────────
|
||||||
POKE_TUNNEL="${POKE_TUNNEL:-1}"
|
|
||||||
|
|
||||||
if [ "${POKE_TUNNEL}" != "1" ]; then
|
if [ "${POKE_TUNNEL}" != "1" ]; then
|
||||||
: "${MCP_API_KEY:?MCP_API_KEY is not set — add it to .env or export it}"
|
: "${MCP_API_KEY:?MCP_API_KEY is not set — add it to .env or export it}"
|
||||||
else
|
else
|
||||||
|
|||||||
Reference in New Issue
Block a user