Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -496,6 +496,7 @@ def create_job(
enabled_toolsets: Optional[List[str]] = None,
workdir: Optional[str] = None,
no_agent: bool = False,
offer_context_inject: bool = False,
) -> Dict[str, Any]:
"""
Create a new cron job.
Expand Down Expand Up @@ -540,6 +541,11 @@ def create_job(
and deliver its stdout directly. Empty stdout = silent (no
delivery). Requires ``script`` to be set. Ideal for classic
watchdogs and periodic alerts that don't need LLM reasoning.
offer_context_inject: When True and deliver targets Telegram, the cron
delivery message includes an inline button that lets the user
manually inject the raw output into the active session context.
If no session is active, tapping the button opens a new session
pre-seeded with the cron output. Defaults to False.

Returns:
The created job dict
Expand Down Expand Up @@ -574,6 +580,7 @@ def create_job(
normalized_toolsets = normalized_toolsets or None
normalized_workdir = _normalize_workdir(workdir)
normalized_no_agent = bool(no_agent)
normalized_offer_context_inject = bool(offer_context_inject)

# no_agent jobs are meaningless without a script β€” the script IS the job.
# Surface this as a clear ValueError at create time so bad configs never
Expand Down Expand Up @@ -627,6 +634,10 @@ def create_job(
"origin": origin, # Tracks where job was created for "origin" delivery
"enabled_toolsets": normalized_toolsets,
"workdir": normalized_workdir,
# Context injection: when True, cron delivery on Telegram includes an
# inline button so the user can manually inject the output into the
# active session context (or start a new session pre-seeded with it).
"offer_context_inject": normalized_offer_context_inject,
}

jobs = load_jobs()
Expand Down
62 changes: 60 additions & 2 deletions cron/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,51 @@

logger = logging.getLogger(__name__)

# ---------------------------------------------------------------------------
# Context-inject store
# ---------------------------------------------------------------------------
# Maps a short token β†’ raw cron output so the Telegram callback handler can
# retrieve the content when the user taps "πŸ“Œ Tambahkan ke sesi" or
# "πŸ’¬ Mulai sesi baru". Entries are keyed by a random hex token embedded in
# the callback_data string ("ci:<token>") and expire after 24 h to avoid
# unbounded memory growth.
# ---------------------------------------------------------------------------
import threading as _threading
import time as _time
import secrets as _secrets

_ci_store: dict = {} # token β†’ {"content": str, "job_id": str, "ts": float}
_ci_store_lock = _threading.Lock()
_CI_TTL = 86400 # 24 hours


def _ci_store_put(content: str, job_id: str) -> str:
"""Store raw cron output and return a short opaque token."""
token = _secrets.token_hex(8)
now = _time.monotonic()
with _ci_store_lock:
# Evict expired entries opportunistically
expired = [k for k, v in _ci_store.items() if now - v["ts"] > _CI_TTL]
for k in expired:
del _ci_store[k]
_ci_store[token] = {"content": content, "job_id": job_id, "ts": now}
return token


def ci_store_pop(token: str) -> Optional[dict]:
"""Retrieve and remove a stored context-inject entry by token.

Returns the stored dict (keys: ``content``, ``job_id``) or None if the
token is unknown or expired. Called by the Telegram callback handler.
"""
with _ci_store_lock:
entry = _ci_store.pop(token, None)
if entry is None:
return None
if _time.monotonic() - entry["ts"] > _CI_TTL:
return None
return entry


class CronPromptInjectionBlocked(Exception):
"""Raised by _build_job_prompt when the fully-assembled prompt trips the
Expand Down Expand Up @@ -578,8 +623,18 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option
# rooms (e.g. Matrix) where the standalone HTTP path cannot encrypt.
runtime_adapter = (adapters or {}).get(platform)
delivered = False

if runtime_adapter is not None and loop is not None and getattr(loop, "is_running", lambda: False)():
send_metadata = {"thread_id": thread_id} if thread_id else None
# Build optional context-inject inline keyboard for Telegram deliveries.
# Only store the token when we know the live adapter will actually send it.
send_metadata = {"thread_id": thread_id} if thread_id else {}
if (
job.get("offer_context_inject")
and platform_name.lower() == "telegram"
):
ci_token = _ci_store_put(content, job.get("id", ""))
send_metadata["cron_context_inject_token"] = ci_token
send_metadata = send_metadata or None
try:
# Send cleaned text (MEDIA tags stripped) β€” not the raw content
text_to_send = cleaned_delivery_content.strip()
Expand Down Expand Up @@ -624,7 +679,10 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option
)

if not delivered:
# Standalone path: run the async send in a fresh event loop (safe from any thread)
# Standalone path: run the async send in a fresh event loop (safe from any thread).
# Note: context-inject button (offer_context_inject) is only supported via the
# live adapter path above. When the gateway is not running, the button is skipped
# and the message is delivered as plain text.
coro = _send_to_platform(platform, pconfig, chat_id, cleaned_delivery_content, thread_id=thread_id, media_files=media_files)
try:
result = asyncio.run(coro)
Expand Down
18 changes: 16 additions & 2 deletions gateway/mirror.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,20 +22,34 @@
_SESSIONS_INDEX = _SESSIONS_DIR / "sessions.json"


def has_active_session(
platform: str,
chat_id: str,
thread_id: Optional[str] = None,
) -> bool:
"""Return True if an active session exists for the given platform + chat_id."""
return _find_session_id(platform, str(chat_id), thread_id=thread_id) is not None


def mirror_to_session(
platform: str,
chat_id: str,
message_text: str,
source_label: str = "cli",
thread_id: Optional[str] = None,
user_id: Optional[str] = None,
role: str = "assistant",
) -> bool:
"""
Append a delivery-mirror message to the target session's transcript.

Finds the gateway session that matches the given platform + chat_id,
then writes a mirror entry to both the JSONL transcript and SQLite DB.

Args:
role: Message role β€” "assistant" for outbound delivery mirrors,
"user" for context-inject (cron output injected as user input).

Returns True if mirrored successfully, False if no matching session or error.
All errors are caught -- this is never fatal.
"""
Expand All @@ -57,7 +71,7 @@ def mirror_to_session(
return False

mirror_msg = {
"role": "assistant",
"role": role,
"content": message_text,
"timestamp": datetime.now().isoformat(),
"mirror": True,
Expand All @@ -67,7 +81,7 @@ def mirror_to_session(
_append_to_jsonl(session_id, mirror_msg)
_append_to_sqlite(session_id, mirror_msg)

logger.debug("Mirror: wrote to session %s (from %s)", session_id, source_label)
logger.debug("Mirror: wrote to session %s (from %s, role=%s)", session_id, source_label, role)
return True

except Exception as e:
Expand Down
146 changes: 145 additions & 1 deletion gateway/platforms/telegram.py
Original file line number Diff line number Diff line change
Expand Up @@ -1527,7 +1527,11 @@ async def send(

message_ids = []
thread_id = self._metadata_thread_id(metadata)


# Context-inject button: attach to the last chunk of a cron delivery
# when the job has offer_context_inject=True.
ci_token = (metadata or {}).get("cron_context_inject_token")

try:
from telegram.error import NetworkError as _NetErr
except ImportError:
Expand All @@ -1544,6 +1548,7 @@ async def send(
_TimedOut = None # type: ignore[assignment,misc]

for i, chunk in enumerate(chunks):
is_last_chunk = (i == len(chunks) - 1)
metadata_reply_to = self._metadata_reply_to_message_id(metadata)
reply_to_source = reply_to or (
str(metadata_reply_to)
Expand All @@ -1563,6 +1568,19 @@ async def send(
effective_thread_id = thread_kwargs.get("message_thread_id")

msg = None
# Build context-inject keyboard once per chunk (outside retry loop)
ci_keyboard = None
if ci_token and is_last_chunk:
ci_keyboard = InlineKeyboardMarkup([[
InlineKeyboardButton(
"πŸ“Œ Add to session",
callback_data=f"ci:inject:{ci_token}",
),
InlineKeyboardButton(
"πŸ’¬ New session",
callback_data=f"ci:new:{ci_token}",
),
]])
for _send_attempt in range(3):
try:
# Try Markdown first, fall back to plain text if it fails
Expand All @@ -1572,6 +1590,7 @@ async def send(
text=chunk,
parse_mode=ParseMode.MARKDOWN_V2,
reply_to_message_id=reply_to_id,
reply_markup=ci_keyboard,
**thread_kwargs,
**self._link_preview_kwargs(),
**self._notification_kwargs(metadata),
Expand All @@ -1586,6 +1605,7 @@ async def send(
text=plain_chunk,
parse_mode=None,
reply_to_message_id=reply_to_id,
reply_markup=ci_keyboard,
**thread_kwargs,
**self._link_preview_kwargs(),
**self._notification_kwargs(metadata),
Expand Down Expand Up @@ -2897,6 +2917,130 @@ async def _handle_callback_query(
)
return

# --- Cron context-inject callbacks (ci:action:token) ---
if data.startswith("ci:"):
parts = data.split(":", 2)
if len(parts) != 3:
await query.answer(text="Invalid context-inject data.")
return

action = parts[1] # "inject" or "new"
ci_token = parts[2]

caller_id = str(getattr(query.from_user, "id", ""))
if not self._is_callback_user_authorized(
caller_id,
chat_id=query_chat_id,
chat_type=str(query_chat_type) if query_chat_type is not None else None,
thread_id=str(query_thread_id) if query_thread_id is not None else None,
user_name=query_user_name,
):
await query.answer(text="β›” You are not authorized.")
return

# Retrieve and consume the stored cron output
try:
from cron.scheduler import ci_store_pop
entry = ci_store_pop(ci_token)
except Exception as exc:
logger.warning("[%s] ci_store_pop failed: %s", self.name, exc)
entry = None

if not entry:
await query.answer(text="⏰ Token expired or not found.")
try:
await query.edit_message_reply_markup(reply_markup=None)
except Exception:
pass
return

raw_content = entry["content"]
job_id = entry.get("job_id", "")
chat_id_str = str(query_chat_id) if query_chat_id is not None else ""
thread_id_str = str(query_thread_id) if query_thread_id is not None else None
label = "⚠️ Gagal" # default, overwritten below

if action == "inject":
# Inject into the active session for this chat as a user message
# so the agent sees it as new input and can respond.
try:
from gateway.mirror import mirror_to_session, has_active_session
if not has_active_session(
platform="telegram",
chat_id=chat_id_str,
thread_id=thread_id_str,
):
ok = False
else:
ok = mirror_to_session(
platform="telegram",
chat_id=chat_id_str,
message_text=raw_content,
source_label=f"cron:{job_id}",
thread_id=thread_id_str,
role="user",
)
except Exception as exc:
logger.warning("[%s] mirror_to_session failed: %s", self.name, exc)
ok = False

if ok:
await query.answer(text="βœ… Added to active session.")
label = "βœ… Added to session"
else:
# No active session found β€” fall back to "new session" behavior
action = "new"

if action == "new":
# Queue the cron output as a pending user message. The next time
# the user sends a message, the gateway will pick it up and the
# agent will see it as context. We cannot create a brand-new
# session from a callback handler (no GatewaySession access), so
# we write it as a user-role mirror entry β€” it will be visible in
# the next session that opens for this chat.
try:
from gateway.mirror import mirror_to_session
ok = mirror_to_session(
platform="telegram",
chat_id=chat_id_str,
message_text=raw_content,
source_label=f"cron:{job_id}",
thread_id=thread_id_str,
role="user",
)
except Exception as exc:
logger.warning("[%s] new-session context inject failed: %s", self.name, exc)
ok = False

if ok:
await query.answer(text="πŸ’¬ Context queued β€” send a message to continue.")
label = "πŸ’¬ Context queued"
else:
await query.answer(text="⚠️ No session found. Start a conversation first.")
label = "⚠️ No active session"

# Edit the message to remove buttons and show status
try:
await query.edit_message_reply_markup(reply_markup=None)
except Exception:
pass
# Send a brief confirmation in the same chat
if query.message and self._bot and ParseMode:
try:
_ci_thread_kwargs = self._thread_kwargs_for_send(
chat_id_str, thread_id_str, {"thread_id": thread_id_str} if thread_id_str else None
)
await self._bot.send_message(
chat_id=int(chat_id_str),
text=self.format_message(f"{label} β€” send your next message to continue."),
parse_mode=ParseMode.MARKDOWN_V2,
**_ci_thread_kwargs,
**self._link_preview_kwargs(),
)
except Exception as exc:
logger.warning("[%s] ci confirmation send failed: %s", self.name, exc)
return

# --- Update prompt callbacks ---
if not data.startswith("update_prompt:"):
return
Expand Down
Loading