Newer
Older
navi-1 / navi / synapse / settings_store.py
"""Per-user Synapse reaction settings + instruction version history (postgres).

There are two standing documents per user, and they answer different questions:

- `instructions` — the *reaction* contract: which events Navi reacts to and what
  it does about them. Read by the dispatcher meta-pass and (through the run's
  system prompt) by every reaction run.
- `dispatcher_instructions` — the *routing* contract: which event types belong to
  which profile, and what may be treated as the same conversation. Read by the
  dispatcher meta-pass only; it never enters a reaction session.

The agent itself may edit either document (self-improvement) — every edit lands in
synapse_instruction_versions with an `edited_by` label and a `doc` discriminant, so
the history doubles as the audit log of the agent touching its own rules.
"""

from datetime import datetime, timezone
from typing import Literal, get_args

from pydantic import BaseModel

PushTarget = Literal["app", "app_synapse", "synapse"]
CompletionNotify = Literal["always", "important", "never"]
InstructionDoc = Literal["reaction", "dispatcher"]
INSTRUCTION_DOCS: tuple[str, ...] = get_args(InstructionDoc)

_VERSIONS_KEEP = 100  # versions kept per (user, doc); older ones are trimmed

# Which settings column each document lives in. The two are versioned
# independently, so a change to the routing rules never buries a reaction edit.
_DOC_COLUMNS: dict[str, str] = {
    "reaction": "instructions",
    "dispatcher": "dispatcher_instructions",
}
assert set(_DOC_COLUMNS) == set(INSTRUCTION_DOCS), "one column per instruction document"


class SynapseSettings(BaseModel):
    user_id: str
    reactions_enabled: bool = False
    push_target: PushTarget = "app"
    completion_notify: CompletionNotify = "important"
    instructions: str = ""
    dispatcher_instructions: str = ""
    # Idle window for continuing a reaction thread; 0 = never continue.
    reaction_session_ttl_minutes: int = 1440
    updated_at: datetime | None = None

    def document(self, doc: InstructionDoc) -> str:
        """The text of one of the two instruction documents."""
        return getattr(self, _DOC_COLUMNS[doc])

    def set_document(self, doc: InstructionDoc, content: str) -> None:
        setattr(self, _DOC_COLUMNS[doc], content)


class InstructionVersion(BaseModel):
    doc: InstructionDoc = "reaction"
    content: str
    edited_by: str
    reason: str | None = None
    created_at: datetime


class SynapseSettingsStore:
    def __init__(self, pool):
        self._pool = pool

    async def get(self, user_id: str) -> SynapseSettings:
        row = await self._pool.fetchrow(
            "SELECT user_id, reactions_enabled, push_target, completion_notify, instructions, "
            "dispatcher_instructions, reaction_session_ttl_minutes, updated_at "
            "FROM synapse_settings WHERE user_id = $1",
            user_id,
        )
        if row is None:
            return SynapseSettings(user_id=user_id)
        return SynapseSettings(
            user_id=row["user_id"],
            reactions_enabled=row["reactions_enabled"],
            push_target=row["push_target"],
            completion_notify=row["completion_notify"],
            instructions=row["instructions"],
            dispatcher_instructions=row["dispatcher_instructions"],
            reaction_session_ttl_minutes=row["reaction_session_ttl_minutes"],
            updated_at=row["updated_at"],
        )

    async def save(self, settings: SynapseSettings) -> None:
        now = datetime.now(timezone.utc)
        await self._pool.execute(
            """
            INSERT INTO synapse_settings
                (user_id, reactions_enabled, push_target, completion_notify, instructions,
                 dispatcher_instructions, reaction_session_ttl_minutes, updated_at)
            VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
            ON CONFLICT (user_id) DO UPDATE SET
                reactions_enabled = EXCLUDED.reactions_enabled,
                push_target = EXCLUDED.push_target,
                completion_notify = EXCLUDED.completion_notify,
                instructions = EXCLUDED.instructions,
                dispatcher_instructions = EXCLUDED.dispatcher_instructions,
                reaction_session_ttl_minutes = EXCLUDED.reaction_session_ttl_minutes,
                updated_at = EXCLUDED.updated_at
            """,
            settings.user_id, settings.reactions_enabled, settings.push_target,
            settings.completion_notify, settings.instructions,
            settings.dispatcher_instructions, settings.reaction_session_ttl_minutes, now,
        )

    async def save_instructions(
        self, user_id: str, content: str, edited_by: Literal["user", "navi"],
        reason: str | None = None, doc: InstructionDoc = "reaction",
    ) -> None:
        """Update one instruction document and record the edit in its history.

        One transaction: the version row and the settings row never disagree.
        """
        column = _DOC_COLUMNS[doc]
        now = datetime.now(timezone.utc)
        async with self._pool.acquire() as conn, conn.transaction():
            await conn.execute(
                f"""
                INSERT INTO synapse_settings (user_id, {column}, updated_at)
                VALUES ($1, $2, $3)
                ON CONFLICT (user_id) DO UPDATE SET
                    {column} = EXCLUDED.{column},
                    updated_at = EXCLUDED.updated_at
                """,
                user_id, content, now,
            )
            await conn.execute(
                """
                INSERT INTO synapse_instruction_versions (user_id, doc, content, edited_by, reason, created_at)
                VALUES ($1, $2, $3, $4, $5, $6)
                """,
                user_id, doc, content, edited_by, reason, now,
            )
            await conn.execute(
                """
                DELETE FROM synapse_instruction_versions
                WHERE user_id = $1 AND doc = $2 AND id NOT IN (
                    SELECT id FROM synapse_instruction_versions
                    WHERE user_id = $1 AND doc = $2 ORDER BY created_at DESC, id DESC LIMIT $3
                )
                """,
                user_id, doc, _VERSIONS_KEEP,
            )

    async def list_versions(
        self, user_id: str, limit: int = 20, doc: InstructionDoc = "reaction",
    ) -> list[InstructionVersion]:
        rows = await self._pool.fetch(
            "SELECT doc, content, edited_by, reason, created_at "
            "FROM synapse_instruction_versions WHERE user_id = $1 AND doc = $2 "
            "ORDER BY created_at DESC, id DESC LIMIT $3",
            user_id, doc, limit,
        )
        return [InstructionVersion(**r) for r in rows]