Compare commits
14
Commits
53684e6c7c
..
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0e1097bdff | ||
|
|
9adec1028c | ||
|
|
01a9eef06c | ||
|
|
43add3b0fe | ||
|
|
1877efdbde | ||
|
|
773d0874d5 | ||
|
|
efc0ad9451 | ||
|
|
207cdb0ef5 | ||
|
|
b5ed0f73b9 | ||
|
|
9a3492c01a | ||
|
|
46365468fe | ||
|
|
0284ab800a | ||
|
|
9b36888b3a | ||
|
|
1a465d10dd |
@@ -60,3 +60,12 @@ STATIC_CRON_TOKEN=REPLACE_WITH_RANDOM_HEX
|
|||||||
# --- SimpleFIN Bridge (optional — simpler US bank automation, ~$1.50/mo) ---
|
# --- SimpleFIN Bridge (optional — simpler US bank automation, ~$1.50/mo) ---
|
||||||
# Connect accounts at https://bridge.simplefin.org, then paste the access URL below
|
# Connect accounts at https://bridge.simplefin.org, then paste the access URL below
|
||||||
# SIMPLEFIN_URL=
|
# SIMPLEFIN_URL=
|
||||||
|
|
||||||
|
# --- Accounting Agent ---
|
||||||
|
# Firefly API base URL as seen from inside the compose network
|
||||||
|
FIREFLY_API_URL=http://app:8080
|
||||||
|
# Personal Access Token: Firefly UI -> Options -> Profile -> OAuth -> Personal Access Tokens
|
||||||
|
AGENT_FIREFLY_TOKEN=REPLACE_WITH_PERSONAL_ACCESS_TOKEN
|
||||||
|
# Name of the agent's own database inside the existing Postgres container
|
||||||
|
AGENT_DB_NAME=agentdb
|
||||||
|
HEARTBEAT_INTERVAL_MINUTES=60
|
||||||
|
|||||||
@@ -1,2 +1,6 @@
|
|||||||
.env
|
.env
|
||||||
*.log
|
*.log
|
||||||
|
.venv/
|
||||||
|
__pycache__/
|
||||||
|
*.egg-info/
|
||||||
|
.pytest_cache/
|
||||||
|
|||||||
@@ -1 +1,3 @@
|
|||||||
# Firefly III - jbnel.dev Finance Stack
|
# Firefly III - jbnel.dev Finance Stack
|
||||||
|
|
||||||
|
`agent/` — personal accounting agent worker; see `agent/README.md` and `docs/superpowers/specs/`.
|
||||||
|
|||||||
@@ -0,0 +1,4 @@
|
|||||||
|
tests/
|
||||||
|
__pycache__/
|
||||||
|
*.egg-info/
|
||||||
|
.pytest_cache/
|
||||||
@@ -0,0 +1,13 @@
|
|||||||
|
FROM python:3.12-slim
|
||||||
|
|
||||||
|
WORKDIR /app
|
||||||
|
|
||||||
|
COPY pyproject.toml ./
|
||||||
|
COPY src ./src
|
||||||
|
|
||||||
|
RUN pip install --no-cache-dir .
|
||||||
|
|
||||||
|
HEALTHCHECK --interval=60s --timeout=10s --start-period=30s \
|
||||||
|
CMD ["python", "-m", "agent.healthcheck"]
|
||||||
|
|
||||||
|
CMD ["python", "-m", "agent"]
|
||||||
@@ -0,0 +1,40 @@
|
|||||||
|
# firefly-agent
|
||||||
|
|
||||||
|
Worker container for the personal accounting agent. See the architecture spec:
|
||||||
|
`docs/superpowers/specs/2026-08-23-accounting-agent-architecture-design.md`.
|
||||||
|
|
||||||
|
## What it does (Phase 1)
|
||||||
|
|
||||||
|
On startup: creates `agentdb` in the stack's Postgres if missing, applies SQL
|
||||||
|
migrations from `src/agent/migrations/`, runs one heartbeat, then schedules the
|
||||||
|
heartbeat every `HEARTBEAT_INTERVAL_MINUTES`. The heartbeat calls Firefly's
|
||||||
|
`/api/v1/about` and records the outcome in `sync_runs`.
|
||||||
|
|
||||||
|
## Development
|
||||||
|
|
||||||
|
python3 -m venv .venv && source .venv/bin/activate
|
||||||
|
pip install -e './agent[dev]'
|
||||||
|
cd agent && pytest # needs Docker running (testcontainers)
|
||||||
|
|
||||||
|
## Configuration
|
||||||
|
|
||||||
|
All via environment variables (see `.env.example`, "Accounting Agent" section):
|
||||||
|
`FIREFLY_API_URL`, `AGENT_FIREFLY_TOKEN` (required), `AGENT_DB_NAME`,
|
||||||
|
`HEARTBEAT_INTERVAL_MINUTES`, plus the existing `DB_HOST`/`DB_PORT`/
|
||||||
|
`POSTGRES_USER`/`POSTGRES_PASSWORD`.
|
||||||
|
|
||||||
|
## Deploy & verify
|
||||||
|
|
||||||
|
# on the host, in the repo directory
|
||||||
|
git pull
|
||||||
|
docker compose build agent
|
||||||
|
docker compose up -d agent
|
||||||
|
docker compose exec agent python -m agent.verify
|
||||||
|
|
||||||
|
Expected verify output: agentdb OK with migrations applied, recent heartbeat
|
||||||
|
runs, Firefly version, and your asset accounts listed.
|
||||||
|
|
||||||
|
Inspect heartbeats directly:
|
||||||
|
|
||||||
|
docker compose exec db psql -U firefly -d agentdb \
|
||||||
|
-c "SELECT connector, status, finished_at FROM sync_runs ORDER BY id DESC LIMIT 5;"
|
||||||
@@ -0,0 +1,26 @@
|
|||||||
|
[project]
|
||||||
|
name = "firefly-agent"
|
||||||
|
version = "0.1.0"
|
||||||
|
description = "Personal accounting agent worker for the Firefly III stack"
|
||||||
|
requires-python = ">=3.12"
|
||||||
|
dependencies = [
|
||||||
|
"apscheduler>=3.10,<4",
|
||||||
|
"httpx>=0.27",
|
||||||
|
"psycopg[binary]>=3.1",
|
||||||
|
]
|
||||||
|
|
||||||
|
[project.optional-dependencies]
|
||||||
|
dev = [
|
||||||
|
"pytest>=8",
|
||||||
|
"testcontainers[postgres]>=4",
|
||||||
|
]
|
||||||
|
|
||||||
|
[build-system]
|
||||||
|
requires = ["hatchling"]
|
||||||
|
build-backend = "hatchling.build"
|
||||||
|
|
||||||
|
[tool.hatch.build.targets.wheel]
|
||||||
|
packages = ["src/agent"]
|
||||||
|
|
||||||
|
[tool.pytest.ini_options]
|
||||||
|
testpaths = ["tests"]
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
__version__ = "0.1.0"
|
||||||
@@ -0,0 +1,37 @@
|
|||||||
|
"""Agent entrypoint: bootstrap the database, then run scheduled jobs forever."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
|
from agent import __version__
|
||||||
|
from agent.config import load_config
|
||||||
|
from agent.db import ensure_database, run_migrations
|
||||||
|
from agent.jobs.heartbeat import run_heartbeat
|
||||||
|
from agent.scheduler import build_scheduler
|
||||||
|
|
||||||
|
logging.basicConfig(
|
||||||
|
level=logging.INFO,
|
||||||
|
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
|
||||||
|
)
|
||||||
|
logger = logging.getLogger("agent")
|
||||||
|
|
||||||
|
|
||||||
|
def main() -> None:
|
||||||
|
logger.info("firefly-agent %s starting", __version__)
|
||||||
|
cfg = load_config()
|
||||||
|
ensure_database(cfg)
|
||||||
|
applied = run_migrations(cfg)
|
||||||
|
if applied:
|
||||||
|
logger.info("applied migrations: %s", ", ".join(applied))
|
||||||
|
status = run_heartbeat(cfg)
|
||||||
|
logger.info("startup heartbeat: %s", status)
|
||||||
|
logger.info(
|
||||||
|
"scheduler starting; heartbeat every %d minutes",
|
||||||
|
cfg.heartbeat_interval_minutes,
|
||||||
|
)
|
||||||
|
build_scheduler(cfg).start()
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
main()
|
||||||
@@ -0,0 +1,53 @@
|
|||||||
|
"""Agent configuration loaded from environment variables."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
from collections.abc import Mapping
|
||||||
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
|
||||||
|
class ConfigError(ValueError):
|
||||||
|
"""A required environment variable is missing or invalid."""
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class AgentConfig:
|
||||||
|
db_host: str
|
||||||
|
db_port: int
|
||||||
|
db_user: str
|
||||||
|
db_password: str
|
||||||
|
agent_db_name: str
|
||||||
|
firefly_url: str
|
||||||
|
firefly_token: str
|
||||||
|
heartbeat_interval_minutes: int
|
||||||
|
|
||||||
|
|
||||||
|
def _require(env: Mapping[str, str], name: str) -> str:
|
||||||
|
value = env.get(name, "").strip()
|
||||||
|
if not value:
|
||||||
|
raise ConfigError(f"missing required environment variable: {name}")
|
||||||
|
return value
|
||||||
|
|
||||||
|
|
||||||
|
def _int(env: Mapping[str, str], name: str, default: str) -> int:
|
||||||
|
"""Parse an integer environment variable, raising ConfigError if invalid."""
|
||||||
|
try:
|
||||||
|
return int(env.get(name, default))
|
||||||
|
except ValueError:
|
||||||
|
raise ConfigError(f"invalid integer for environment variable: {name}")
|
||||||
|
|
||||||
|
|
||||||
|
def load_config(env: Mapping[str, str] | None = None) -> AgentConfig:
|
||||||
|
if env is None:
|
||||||
|
env = os.environ
|
||||||
|
return AgentConfig(
|
||||||
|
db_host=_require(env, "DB_HOST"),
|
||||||
|
db_port=_int(env, "DB_PORT", "5432"),
|
||||||
|
db_user=_require(env, "POSTGRES_USER"),
|
||||||
|
db_password=_require(env, "POSTGRES_PASSWORD"),
|
||||||
|
agent_db_name=env.get("AGENT_DB_NAME", "agentdb"),
|
||||||
|
firefly_url=_require(env, "FIREFLY_API_URL"),
|
||||||
|
firefly_token=_require(env, "AGENT_FIREFLY_TOKEN"),
|
||||||
|
heartbeat_interval_minutes=_int(env, "HEARTBEAT_INTERVAL_MINUTES", "60"),
|
||||||
|
)
|
||||||
@@ -0,0 +1,74 @@
|
|||||||
|
"""agentdb access: bootstrap, migrations, connections."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import psycopg
|
||||||
|
|
||||||
|
from agent.config import AgentConfig
|
||||||
|
|
||||||
|
_MIGRATIONS_DIR = Path(__file__).parent / "migrations"
|
||||||
|
|
||||||
|
|
||||||
|
def _dsn(cfg: AgentConfig, dbname: str) -> str:
|
||||||
|
from psycopg import conninfo
|
||||||
|
return conninfo.make_conninfo(
|
||||||
|
host=cfg.db_host,
|
||||||
|
port=cfg.db_port,
|
||||||
|
dbname=dbname,
|
||||||
|
user=cfg.db_user,
|
||||||
|
password=cfg.db_password,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def admin_dsn(cfg: AgentConfig) -> str:
|
||||||
|
"""DSN for the maintenance database that always exists."""
|
||||||
|
return _dsn(cfg, "postgres")
|
||||||
|
|
||||||
|
|
||||||
|
def agent_dsn(cfg: AgentConfig) -> str:
|
||||||
|
return _dsn(cfg, cfg.agent_db_name)
|
||||||
|
|
||||||
|
|
||||||
|
def connect(cfg: AgentConfig) -> psycopg.Connection:
|
||||||
|
return psycopg.connect(agent_dsn(cfg))
|
||||||
|
|
||||||
|
|
||||||
|
def ensure_database(cfg: AgentConfig) -> None:
|
||||||
|
"""Create the agent database if it does not exist. Idempotent."""
|
||||||
|
with psycopg.connect(admin_dsn(cfg), autocommit=True) as conn:
|
||||||
|
exists = conn.execute(
|
||||||
|
"SELECT 1 FROM pg_database WHERE datname = %s", (cfg.agent_db_name,)
|
||||||
|
).fetchone()
|
||||||
|
if not exists:
|
||||||
|
# CREATE DATABASE cannot be parameterized; the name comes from
|
||||||
|
# our own config, quoted defensively via psycopg's identifier API.
|
||||||
|
from psycopg import sql
|
||||||
|
|
||||||
|
conn.execute(
|
||||||
|
sql.SQL("CREATE DATABASE {}").format(sql.Identifier(cfg.agent_db_name))
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def run_migrations(cfg: AgentConfig) -> list[str]:
|
||||||
|
"""Apply pending migrations in filename order. Returns names applied."""
|
||||||
|
applied: list[str] = []
|
||||||
|
with connect(cfg) as conn:
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE IF NOT EXISTS schema_migrations ("
|
||||||
|
" name TEXT PRIMARY KEY,"
|
||||||
|
" applied_at TIMESTAMPTZ NOT NULL DEFAULT now())"
|
||||||
|
)
|
||||||
|
done = {
|
||||||
|
row[0] for row in conn.execute("SELECT name FROM schema_migrations")
|
||||||
|
}
|
||||||
|
for path in sorted(_MIGRATIONS_DIR.glob("*.sql")):
|
||||||
|
if path.name in done:
|
||||||
|
continue
|
||||||
|
conn.execute(path.read_text())
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO schema_migrations (name) VALUES (%s)", (path.name,)
|
||||||
|
)
|
||||||
|
applied.append(path.name)
|
||||||
|
return applied
|
||||||
@@ -0,0 +1,54 @@
|
|||||||
|
"""Minimal read-only client for the Firefly III REST API."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
|
||||||
|
class FireflyError(Exception):
|
||||||
|
"""Firefly III returned an unexpected response."""
|
||||||
|
|
||||||
|
|
||||||
|
class FireflyClient:
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
base_url: str,
|
||||||
|
token: str,
|
||||||
|
transport: httpx.BaseTransport | None = None,
|
||||||
|
) -> None:
|
||||||
|
self._http = httpx.Client(
|
||||||
|
base_url=base_url,
|
||||||
|
headers={
|
||||||
|
"Authorization": f"Bearer {token}",
|
||||||
|
"Accept": "application/json",
|
||||||
|
},
|
||||||
|
timeout=30.0,
|
||||||
|
transport=transport,
|
||||||
|
)
|
||||||
|
|
||||||
|
def _get(self, path: str, params: dict | None = None) -> dict:
|
||||||
|
try:
|
||||||
|
response = self._http.get(path, params=params)
|
||||||
|
except httpx.HTTPError as exc:
|
||||||
|
raise FireflyError(
|
||||||
|
f"GET {path} failed: {type(exc).__name__}: {exc}"
|
||||||
|
) from exc
|
||||||
|
if response.status_code != 200:
|
||||||
|
raise FireflyError(
|
||||||
|
f"GET {path} -> {response.status_code}: {response.text[:200]}"
|
||||||
|
)
|
||||||
|
return response.json()
|
||||||
|
|
||||||
|
def about(self) -> dict:
|
||||||
|
return self._get("/api/v1/about")["data"]
|
||||||
|
|
||||||
|
def accounts(self, type: str = "asset") -> list[dict]:
|
||||||
|
results: list[dict] = []
|
||||||
|
page = 1
|
||||||
|
while True:
|
||||||
|
body = self._get("/api/v1/accounts", params={"type": type, "page": page})
|
||||||
|
results.extend(body["data"])
|
||||||
|
pagination = body.get("meta", {}).get("pagination", {})
|
||||||
|
if page >= pagination.get("total_pages", 1):
|
||||||
|
return results
|
||||||
|
page += 1
|
||||||
@@ -0,0 +1,23 @@
|
|||||||
|
"""Docker HEALTHCHECK entry: healthy = agentdb answers SELECT 1."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import sys
|
||||||
|
|
||||||
|
from agent.config import load_config
|
||||||
|
from agent.db import connect
|
||||||
|
|
||||||
|
|
||||||
|
def main() -> int:
|
||||||
|
try:
|
||||||
|
cfg = load_config()
|
||||||
|
with connect(cfg) as conn:
|
||||||
|
conn.execute("SELECT 1")
|
||||||
|
return 0
|
||||||
|
except Exception as exc:
|
||||||
|
print(f"unhealthy: {exc}", file=sys.stderr)
|
||||||
|
return 1
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
sys.exit(main())
|
||||||
@@ -0,0 +1,34 @@
|
|||||||
|
"""Heartbeat: prove Firefly API + agentdb connectivity on a schedule."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
|
||||||
|
from agent.config import AgentConfig
|
||||||
|
from agent.db import connect
|
||||||
|
from agent.firefly import FireflyClient
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
def run_heartbeat(cfg: AgentConfig, client: FireflyClient | None = None) -> str:
|
||||||
|
if client is None:
|
||||||
|
client = FireflyClient(cfg.firefly_url, cfg.firefly_token)
|
||||||
|
|
||||||
|
started_at = datetime.now(timezone.utc)
|
||||||
|
try:
|
||||||
|
info = client.about()
|
||||||
|
status, detail = "ok", f"firefly version {info.get('version', 'unknown')}"
|
||||||
|
except Exception as exc: # any failure is a recorded outcome, not a crash
|
||||||
|
status, detail = "error", f"{type(exc).__name__}: {exc}"
|
||||||
|
logger.warning("heartbeat failed: %s", detail)
|
||||||
|
|
||||||
|
finished_at = datetime.now(timezone.utc)
|
||||||
|
with connect(cfg) as conn:
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO sync_runs (connector, started_at, finished_at, status, detail)"
|
||||||
|
" VALUES (%s, %s, %s, %s, %s)",
|
||||||
|
("heartbeat", started_at, finished_at, status, detail),
|
||||||
|
)
|
||||||
|
return status
|
||||||
@@ -0,0 +1,47 @@
|
|||||||
|
-- agentdb schema v1 (see spec: agentdb schema section)
|
||||||
|
|
||||||
|
CREATE TABLE sync_runs (
|
||||||
|
id BIGSERIAL PRIMARY KEY,
|
||||||
|
connector TEXT NOT NULL,
|
||||||
|
started_at TIMESTAMPTZ NOT NULL,
|
||||||
|
finished_at TIMESTAMPTZ,
|
||||||
|
status TEXT NOT NULL CHECK (status IN ('running', 'ok', 'error')),
|
||||||
|
items_upserted INTEGER NOT NULL DEFAULT 0,
|
||||||
|
detail TEXT
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE INDEX sync_runs_connector_id_idx ON sync_runs (connector, id DESC);
|
||||||
|
|
||||||
|
CREATE TABLE documents (
|
||||||
|
id BIGSERIAL PRIMARY KEY,
|
||||||
|
path TEXT NOT NULL UNIQUE,
|
||||||
|
sha256 CHAR(64) NOT NULL UNIQUE,
|
||||||
|
institution TEXT NOT NULL,
|
||||||
|
account TEXT,
|
||||||
|
doc_type TEXT NOT NULL,
|
||||||
|
period TEXT,
|
||||||
|
source TEXT NOT NULL,
|
||||||
|
received_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE TABLE bills (
|
||||||
|
id BIGSERIAL PRIMARY KEY,
|
||||||
|
payee TEXT NOT NULL,
|
||||||
|
expected_amount NUMERIC(12,2),
|
||||||
|
due_date DATE NOT NULL,
|
||||||
|
funding_account TEXT,
|
||||||
|
autopay BOOLEAN NOT NULL DEFAULT false,
|
||||||
|
status TEXT NOT NULL DEFAULT 'upcoming'
|
||||||
|
CHECK (status IN ('upcoming', 'approved', 'paid', 'skipped', 'overdue'))
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE TABLE approvals (
|
||||||
|
id BIGSERIAL PRIMARY KEY,
|
||||||
|
bill_id BIGINT NOT NULL REFERENCES bills (id),
|
||||||
|
amount NUMERIC(12,2) NOT NULL,
|
||||||
|
requested_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
decided_at TIMESTAMPTZ,
|
||||||
|
decision TEXT CHECK (decision IN ('approved', 'denied')),
|
||||||
|
telegram_message_id TEXT,
|
||||||
|
executed_at TIMESTAMPTZ
|
||||||
|
);
|
||||||
@@ -0,0 +1,23 @@
|
|||||||
|
"""APScheduler assembly: one place where jobs get registered."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from apscheduler.schedulers.base import BaseScheduler
|
||||||
|
from apscheduler.schedulers.blocking import BlockingScheduler
|
||||||
|
|
||||||
|
from agent.config import AgentConfig
|
||||||
|
from agent.jobs.heartbeat import run_heartbeat
|
||||||
|
|
||||||
|
|
||||||
|
def build_scheduler(
|
||||||
|
cfg: AgentConfig, *, scheduler_class: type = BlockingScheduler
|
||||||
|
) -> BaseScheduler:
|
||||||
|
scheduler = scheduler_class(timezone="UTC")
|
||||||
|
scheduler.add_job(
|
||||||
|
run_heartbeat,
|
||||||
|
"interval",
|
||||||
|
args=[cfg],
|
||||||
|
minutes=cfg.heartbeat_interval_minutes,
|
||||||
|
id="heartbeat",
|
||||||
|
)
|
||||||
|
return scheduler
|
||||||
@@ -0,0 +1,42 @@
|
|||||||
|
"""Post-deploy verification: `python -m agent.verify`."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from agent.config import AgentConfig, load_config
|
||||||
|
from agent.db import connect
|
||||||
|
from agent.firefly import FireflyClient
|
||||||
|
|
||||||
|
|
||||||
|
def verify(cfg: AgentConfig, client: FireflyClient | None = None) -> list[str]:
|
||||||
|
lines: list[str] = []
|
||||||
|
|
||||||
|
with connect(cfg) as conn:
|
||||||
|
migrations = conn.execute("SELECT count(*) FROM schema_migrations").fetchone()[0]
|
||||||
|
runs = conn.execute(
|
||||||
|
"SELECT connector, status, finished_at"
|
||||||
|
" FROM sync_runs ORDER BY id DESC LIMIT 3"
|
||||||
|
).fetchall()
|
||||||
|
lines.append(f"agentdb: OK ({migrations} migration(s) applied)")
|
||||||
|
for connector, status, finished_at in runs:
|
||||||
|
lines.append(f" last run: {connector} {status} at {finished_at}")
|
||||||
|
|
||||||
|
if client is None:
|
||||||
|
client = FireflyClient(cfg.firefly_url, cfg.firefly_token)
|
||||||
|
info = client.about()
|
||||||
|
lines.append(f"firefly: OK (version {info.get('version', 'unknown')})")
|
||||||
|
|
||||||
|
accounts = client.accounts()
|
||||||
|
lines.append(f"firefly: {len(accounts)} asset account(s):")
|
||||||
|
for account in accounts:
|
||||||
|
lines.append(f" - {account['attributes']['name']}")
|
||||||
|
return lines
|
||||||
|
|
||||||
|
|
||||||
|
def main() -> int:
|
||||||
|
for line in verify(load_config()):
|
||||||
|
print(line)
|
||||||
|
return 0
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
raise SystemExit(main())
|
||||||
@@ -0,0 +1,29 @@
|
|||||||
|
import itertools
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from testcontainers.community.postgres import PostgresContainer
|
||||||
|
|
||||||
|
from agent.config import AgentConfig
|
||||||
|
|
||||||
|
_db_counter = itertools.count()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(scope="session")
|
||||||
|
def pg():
|
||||||
|
with PostgresContainer("postgres:16-alpine") as container:
|
||||||
|
yield container
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def cfg(pg) -> AgentConfig:
|
||||||
|
"""AgentConfig pointing at a fresh, uniquely named database per test."""
|
||||||
|
return AgentConfig(
|
||||||
|
db_host=pg.get_container_host_ip(),
|
||||||
|
db_port=int(pg.get_exposed_port(5432)),
|
||||||
|
db_user=pg.username,
|
||||||
|
db_password=pg.password,
|
||||||
|
agent_db_name=f"agentdb_test_{next(_db_counter)}",
|
||||||
|
firefly_url="http://firefly.test",
|
||||||
|
firefly_token="test-token",
|
||||||
|
heartbeat_interval_minutes=60,
|
||||||
|
)
|
||||||
@@ -0,0 +1,65 @@
|
|||||||
|
import pytest
|
||||||
|
|
||||||
|
from agent.config import AgentConfig, ConfigError, load_config
|
||||||
|
|
||||||
|
FULL_ENV = {
|
||||||
|
"DB_HOST": "db",
|
||||||
|
"DB_PORT": "5433",
|
||||||
|
"POSTGRES_USER": "firefly",
|
||||||
|
"POSTGRES_PASSWORD": "s3cret",
|
||||||
|
"AGENT_DB_NAME": "agentdb2",
|
||||||
|
"FIREFLY_API_URL": "http://app:8080",
|
||||||
|
"AGENT_FIREFLY_TOKEN": "tok123",
|
||||||
|
"HEARTBEAT_INTERVAL_MINUTES": "15",
|
||||||
|
}
|
||||||
|
|
||||||
|
MINIMAL_ENV = {
|
||||||
|
"DB_HOST": "db",
|
||||||
|
"POSTGRES_USER": "firefly",
|
||||||
|
"POSTGRES_PASSWORD": "s3cret",
|
||||||
|
"FIREFLY_API_URL": "http://app:8080",
|
||||||
|
"AGENT_FIREFLY_TOKEN": "tok123",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_load_config_reads_all_vars():
|
||||||
|
cfg = load_config(FULL_ENV)
|
||||||
|
assert cfg == AgentConfig(
|
||||||
|
db_host="db",
|
||||||
|
db_port=5433,
|
||||||
|
db_user="firefly",
|
||||||
|
db_password="s3cret",
|
||||||
|
agent_db_name="agentdb2",
|
||||||
|
firefly_url="http://app:8080",
|
||||||
|
firefly_token="tok123",
|
||||||
|
heartbeat_interval_minutes=15,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_defaults_applied():
|
||||||
|
cfg = load_config(MINIMAL_ENV)
|
||||||
|
assert cfg.db_port == 5432
|
||||||
|
assert cfg.agent_db_name == "agentdb"
|
||||||
|
assert cfg.heartbeat_interval_minutes == 60
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"missing",
|
||||||
|
["DB_HOST", "POSTGRES_USER", "POSTGRES_PASSWORD", "FIREFLY_API_URL", "AGENT_FIREFLY_TOKEN"],
|
||||||
|
)
|
||||||
|
def test_missing_required_var_names_the_var(missing):
|
||||||
|
env = {k: v for k, v in MINIMAL_ENV.items() if k != missing}
|
||||||
|
with pytest.raises(ConfigError, match=missing):
|
||||||
|
load_config(env)
|
||||||
|
|
||||||
|
|
||||||
|
def test_invalid_db_port_raises_config_error():
|
||||||
|
env = {**MINIMAL_ENV, "DB_PORT": "notanumber"}
|
||||||
|
with pytest.raises(ConfigError, match="DB_PORT"):
|
||||||
|
load_config(env)
|
||||||
|
|
||||||
|
|
||||||
|
def test_invalid_heartbeat_interval_raises_config_error():
|
||||||
|
env = {**MINIMAL_ENV, "HEARTBEAT_INTERVAL_MINUTES": "abc"}
|
||||||
|
with pytest.raises(ConfigError, match="HEARTBEAT_INTERVAL_MINUTES"):
|
||||||
|
load_config(env)
|
||||||
@@ -0,0 +1,61 @@
|
|||||||
|
from agent.db import connect, ensure_database, run_migrations
|
||||||
|
|
||||||
|
EXPECTED_TABLES = {"schema_migrations", "sync_runs", "documents", "bills", "approvals"}
|
||||||
|
|
||||||
|
|
||||||
|
def test_ensure_database_creates_and_is_idempotent(cfg):
|
||||||
|
ensure_database(cfg)
|
||||||
|
ensure_database(cfg) # second call must not raise
|
||||||
|
with connect(cfg) as conn:
|
||||||
|
assert conn.execute("SELECT 1").fetchone()[0] == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_run_migrations_applies_each_migration_exactly_once(cfg):
|
||||||
|
ensure_database(cfg)
|
||||||
|
|
||||||
|
first = run_migrations(cfg)
|
||||||
|
assert first == ["0001_initial.sql"]
|
||||||
|
|
||||||
|
second = run_migrations(cfg)
|
||||||
|
assert second == []
|
||||||
|
|
||||||
|
with connect(cfg) as conn:
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT tablename FROM pg_tables WHERE schemaname = 'public'"
|
||||||
|
).fetchall()
|
||||||
|
assert EXPECTED_TABLES <= {r[0] for r in rows}
|
||||||
|
|
||||||
|
|
||||||
|
def test_sync_runs_rejects_bad_status(cfg):
|
||||||
|
import psycopg
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
ensure_database(cfg)
|
||||||
|
run_migrations(cfg)
|
||||||
|
with connect(cfg) as conn, pytest.raises(psycopg.errors.CheckViolation):
|
||||||
|
conn.execute(
|
||||||
|
"INSERT INTO sync_runs (connector, started_at, status)"
|
||||||
|
" VALUES ('x', now(), 'bogus')"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_agent_dsn_escapes_password_with_space_and_quote():
|
||||||
|
import psycopg
|
||||||
|
from agent.config import AgentConfig
|
||||||
|
from agent.db import agent_dsn
|
||||||
|
|
||||||
|
cfg = AgentConfig(
|
||||||
|
db_host="localhost",
|
||||||
|
db_port=5432,
|
||||||
|
db_user="testuser",
|
||||||
|
db_password="pa ss'word", # password with space and single quote
|
||||||
|
agent_db_name="testdb",
|
||||||
|
firefly_url="http://firefly.test",
|
||||||
|
firefly_token="test-token",
|
||||||
|
heartbeat_interval_minutes=60,
|
||||||
|
)
|
||||||
|
|
||||||
|
dsn = agent_dsn(cfg)
|
||||||
|
# Parse DSN and verify password is preserved correctly
|
||||||
|
parsed = psycopg.conninfo.conninfo_to_dict(dsn)
|
||||||
|
assert parsed["password"] == "pa ss'word"
|
||||||
@@ -0,0 +1,57 @@
|
|||||||
|
import httpx
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from agent.firefly import FireflyClient, FireflyError
|
||||||
|
|
||||||
|
|
||||||
|
def make_client(handler) -> FireflyClient:
|
||||||
|
return FireflyClient(
|
||||||
|
"http://firefly.test", "tok", transport=httpx.MockTransport(handler)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_about_sends_bearer_token_and_returns_data():
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
assert request.headers["Authorization"] == "Bearer tok"
|
||||||
|
assert request.headers["Accept"] == "application/json"
|
||||||
|
assert request.url.path == "/api/v1/about"
|
||||||
|
return httpx.Response(200, json={"data": {"version": "6.1.1"}})
|
||||||
|
|
||||||
|
assert make_client(handler).about() == {"version": "6.1.1"}
|
||||||
|
|
||||||
|
|
||||||
|
def test_non_200_raises_firefly_error_with_status():
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
return httpx.Response(401, json={"message": "Unauthenticated"})
|
||||||
|
|
||||||
|
with pytest.raises(FireflyError, match="401"):
|
||||||
|
make_client(handler).about()
|
||||||
|
|
||||||
|
|
||||||
|
def test_accounts_walks_all_pages():
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
assert request.url.path == "/api/v1/accounts"
|
||||||
|
assert request.url.params["type"] == "asset"
|
||||||
|
page = int(request.url.params.get("page", "1"))
|
||||||
|
data = {
|
||||||
|
1: [{"id": "1", "attributes": {"name": "NFCU Checking"}}],
|
||||||
|
2: [{"id": "2", "attributes": {"name": "Venmo"}}],
|
||||||
|
}[page]
|
||||||
|
return httpx.Response(
|
||||||
|
200,
|
||||||
|
json={
|
||||||
|
"data": data,
|
||||||
|
"meta": {"pagination": {"total_pages": 2, "current_page": page}},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
accounts = make_client(handler).accounts()
|
||||||
|
assert [a["attributes"]["name"] for a in accounts] == ["NFCU Checking", "Venmo"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_transport_errors_wrapped_in_firefly_error():
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
raise httpx.ConnectError("connection refused")
|
||||||
|
|
||||||
|
with pytest.raises(FireflyError, match="ConnectError"):
|
||||||
|
make_client(handler).about()
|
||||||
@@ -0,0 +1,48 @@
|
|||||||
|
from agent.db import connect, ensure_database, run_migrations
|
||||||
|
from agent.firefly import FireflyError
|
||||||
|
from agent.jobs.heartbeat import run_heartbeat
|
||||||
|
|
||||||
|
|
||||||
|
class FakeFirefly:
|
||||||
|
def __init__(self, version: str | None = None, error: Exception | None = None):
|
||||||
|
self._version = version
|
||||||
|
self._error = error
|
||||||
|
|
||||||
|
def about(self) -> dict:
|
||||||
|
if self._error:
|
||||||
|
raise self._error
|
||||||
|
return {"version": self._version}
|
||||||
|
|
||||||
|
|
||||||
|
def _last_run(cfg):
|
||||||
|
with connect(cfg) as conn:
|
||||||
|
return conn.execute(
|
||||||
|
"SELECT connector, status, detail, started_at, finished_at"
|
||||||
|
" FROM sync_runs ORDER BY id DESC LIMIT 1"
|
||||||
|
).fetchone()
|
||||||
|
|
||||||
|
|
||||||
|
def test_heartbeat_records_ok(cfg):
|
||||||
|
ensure_database(cfg)
|
||||||
|
run_migrations(cfg)
|
||||||
|
|
||||||
|
status = run_heartbeat(cfg, client=FakeFirefly(version="6.1.1"))
|
||||||
|
|
||||||
|
assert status == "ok"
|
||||||
|
connector, db_status, detail, started_at, finished_at = _last_run(cfg)
|
||||||
|
assert connector == "heartbeat"
|
||||||
|
assert db_status == "ok"
|
||||||
|
assert "6.1.1" in detail
|
||||||
|
assert finished_at >= started_at
|
||||||
|
|
||||||
|
|
||||||
|
def test_heartbeat_records_error_and_does_not_raise(cfg):
|
||||||
|
ensure_database(cfg)
|
||||||
|
run_migrations(cfg)
|
||||||
|
|
||||||
|
status = run_heartbeat(cfg, client=FakeFirefly(error=FireflyError("GET /api/v1/about -> 401")))
|
||||||
|
|
||||||
|
assert status == "error"
|
||||||
|
connector, db_status, detail, *_ = _last_run(cfg)
|
||||||
|
assert db_status == "error"
|
||||||
|
assert "401" in detail
|
||||||
@@ -0,0 +1,30 @@
|
|||||||
|
from datetime import timedelta
|
||||||
|
|
||||||
|
from apscheduler.schedulers.background import BackgroundScheduler
|
||||||
|
|
||||||
|
from agent.config import AgentConfig
|
||||||
|
from agent.scheduler import build_scheduler
|
||||||
|
|
||||||
|
|
||||||
|
def _cfg(minutes: int) -> AgentConfig:
|
||||||
|
return AgentConfig(
|
||||||
|
db_host="db",
|
||||||
|
db_port=5432,
|
||||||
|
db_user="u",
|
||||||
|
db_password="p",
|
||||||
|
agent_db_name="agentdb",
|
||||||
|
firefly_url="http://app:8080",
|
||||||
|
firefly_token="tok",
|
||||||
|
heartbeat_interval_minutes=minutes,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_scheduler_registers_heartbeat_at_configured_interval():
|
||||||
|
scheduler = build_scheduler(_cfg(minutes=15), scheduler_class=BackgroundScheduler)
|
||||||
|
scheduler.start(paused=True)
|
||||||
|
try:
|
||||||
|
job = scheduler.get_job("heartbeat")
|
||||||
|
assert job is not None
|
||||||
|
assert job.trigger.interval == timedelta(minutes=15)
|
||||||
|
finally:
|
||||||
|
scheduler.shutdown(wait=False)
|
||||||
@@ -0,0 +1,4 @@
|
|||||||
|
def test_package_imports():
|
||||||
|
import agent
|
||||||
|
|
||||||
|
assert agent.__version__ == "0.1.0"
|
||||||
@@ -0,0 +1,28 @@
|
|||||||
|
from agent.db import ensure_database, run_migrations
|
||||||
|
from agent.jobs.heartbeat import run_heartbeat
|
||||||
|
from agent.verify import verify
|
||||||
|
|
||||||
|
|
||||||
|
class FakeFirefly:
|
||||||
|
def about(self) -> dict:
|
||||||
|
return {"version": "6.1.1"}
|
||||||
|
|
||||||
|
def accounts(self, type: str = "asset") -> list[dict]:
|
||||||
|
return [
|
||||||
|
{"id": "1", "attributes": {"name": "NFCU Checking"}},
|
||||||
|
{"id": "2", "attributes": {"name": "Venmo"}},
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def test_verify_reports_db_firefly_and_accounts(cfg):
|
||||||
|
ensure_database(cfg)
|
||||||
|
run_migrations(cfg)
|
||||||
|
run_heartbeat(cfg, client=FakeFirefly())
|
||||||
|
|
||||||
|
report = "\n".join(verify(cfg, client=FakeFirefly()))
|
||||||
|
|
||||||
|
assert "1 migration(s) applied" in report
|
||||||
|
assert "heartbeat ok" in report
|
||||||
|
assert "firefly: OK (version 6.1.1)" in report
|
||||||
|
assert "2 asset account(s)" in report
|
||||||
|
assert "NFCU Checking" in report
|
||||||
@@ -66,6 +66,18 @@ services:
|
|||||||
- "traefik.http.services.firefly-importer.loadbalancer.server.port=8080"
|
- "traefik.http.services.firefly-importer.loadbalancer.server.port=8080"
|
||||||
- "traefik.docker.network=traefik_proxy"
|
- "traefik.docker.network=traefik_proxy"
|
||||||
|
|
||||||
|
agent:
|
||||||
|
build: ./agent
|
||||||
|
restart: unless-stopped
|
||||||
|
env_file: .env
|
||||||
|
networks:
|
||||||
|
- firefly
|
||||||
|
depends_on:
|
||||||
|
db:
|
||||||
|
condition: service_healthy
|
||||||
|
app:
|
||||||
|
condition: service_healthy
|
||||||
|
|
||||||
db:
|
db:
|
||||||
image: postgres:16-alpine
|
image: postgres:16-alpine
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,155 @@
|
|||||||
|
# Personal Accounting Agent — Architecture Design
|
||||||
|
|
||||||
|
**Date:** 2026-08-23
|
||||||
|
**Status:** Approved (Jake, 2026-08-23)
|
||||||
|
**Scope:** Master architecture and phase decomposition. Each phase gets its own spec and implementation plan; this document is the frame they hang on.
|
||||||
|
|
||||||
|
## Purpose
|
||||||
|
|
||||||
|
Build a personal accounting agent around the existing Firefly III stack at
|
||||||
|
`https://firefly.jbnel.dev`. The agent:
|
||||||
|
|
||||||
|
1. Monitors account balances across all institutions.
|
||||||
|
2. Assesses spending trends, tracks budgets, and forecasts cashflow.
|
||||||
|
3. Prepares utility and other bill payments, executing only with Jake's explicit approval.
|
||||||
|
4. Alerts Jake to pertinent information through Telegram, ntfy, and email.
|
||||||
|
5. Maintains account records (statements, bill PDFs) in a structured local folder, with cloud sync later.
|
||||||
|
6. Keeps all record data queryable — ledger data in Firefly, everything else in a dedicated agent database.
|
||||||
|
|
||||||
|
## Decisions
|
||||||
|
|
||||||
|
These were settled in the design conversation and bind all phases:
|
||||||
|
|
||||||
|
| Decision | Choice |
|
||||||
|
|---|---|
|
||||||
|
| Primary bank feed | SimpleFIN Bridge (read-only, ~$1.50/mo) |
|
||||||
|
| Coverage strategy | Connector layer: SimpleFIN plus per-institution connectors until NFCU, TFCU, Cash App, Venmo, and PayPal are all covered. Verify real coverage per institution early; assume nothing. |
|
||||||
|
| Wallet feeds (Cash App, Venmo, PayPal) | *(Amended 2026-08-24 — Jake wants automation, no manual CSV.)* Primary: per-transaction email receipts parsed from the finance inbox (near-real-time; these apps support no aggregators). Backstop: browser automation fetches monthly statements/CSVs for reconciliation. Manual CSV drop folder is break-glass fallback only. PayPal may later upgrade to its Transaction Search API via business account. |
|
||||||
|
| Finance inbox | `finance@ordinatorlabs.com` — dedicated IONOS mailbox (ordinatorlabs.com already has IONOS MX; jbnel.dev has no mail hosting). Agent polls IONOS IMAP. Optional later: forward `finance@jbnel.dev` into it. |
|
||||||
|
| Bill-pay autonomy | Prep + approval. The agent tracks, verifies funding, and asks. Money never moves without an explicit approval from Jake. |
|
||||||
|
| Runtime | Hybrid. One Python worker container does deterministic work on schedules; scheduled Claude sessions do judgment work (analysis, forecast narrative, anomaly triage). |
|
||||||
|
| Notification channels | Telegram bot (interactive: alerts, approvals), ntfy (urgent push), email (digests, receipts). Routed by message class. |
|
||||||
|
| Document acquisition | Dedicated finance email inbox via IMAP first; browser-automation downloads later for billers that never email PDFs. |
|
||||||
|
| Business accounts | Separate Firefly instance when a venture opens. This instance stays purely personal. The connector layer is account-agnostic so it can serve both. |
|
||||||
|
| Agent state | `agentdb`: a separate database in the existing Postgres container. |
|
||||||
|
|
||||||
|
## Architecture
|
||||||
|
|
||||||
|
One new container, `agent`, joins the existing compose stack (Firefly III core,
|
||||||
|
Data Importer, Postgres 16, Redis, cron) behind Traefik. The worker is a Python
|
||||||
|
3.12 application under `agent/` in this repo, structured as modules with hard
|
||||||
|
boundaries so any module can become its own service during the planned k8s
|
||||||
|
migration:
|
||||||
|
|
||||||
|
- **Connectors** — pull data from sources on schedules. Each connector
|
||||||
|
implements one interface: fetch → emit canonical records. Initial set:
|
||||||
|
SimpleFIN poller, CSV drop-folder watcher, IMAP inbox poller. Later: browser
|
||||||
|
automation jobs.
|
||||||
|
- **Sync engine** — normalizes connector output into canonical models
|
||||||
|
(transaction, balance, document, bill), deduplicates, and pushes to Firefly
|
||||||
|
via its REST API using a personal access token.
|
||||||
|
- **Notifier** — one interface, three sinks (Telegram, ntfy, SMTP). Message
|
||||||
|
class determines routing: interactive requests and real-time alerts go to
|
||||||
|
Telegram, urgent one-way pushes to ntfy, digests and receipts to email.
|
||||||
|
- **Rules engine** — deterministic checks after each sync: low balance, large
|
||||||
|
or unusual transaction, failed sync, upcoming bill without positioned funds.
|
||||||
|
- **Bill module** — maintains the bill calendar in `agentdb` from Firefly
|
||||||
|
bills and parsed bill emails; runs the approval workflow.
|
||||||
|
- **Scheduler** — APScheduler; every job idempotent and individually
|
||||||
|
triggerable for testing.
|
||||||
|
|
||||||
|
Claude sessions run on schedules outside the container (Claude Code
|
||||||
|
routines/cron). They read the Firefly API and `agentdb`, write analysis and
|
||||||
|
narratives out through the notifier, and never write directly to the ledger.
|
||||||
|
The agent container exposes the notifier to them as one token-authenticated
|
||||||
|
HTTP endpoint (`POST /notify`), reachable only on the internal Docker network
|
||||||
|
or via SSH tunnel — never through Traefik.
|
||||||
|
|
||||||
|
### Data flow
|
||||||
|
|
||||||
|
```
|
||||||
|
SimpleFIN ─┐
|
||||||
|
CSV drop ─┼→ Connectors → canonical models → dedup → Firefly API (ledger)
|
||||||
|
IMAP ─┤ │
|
||||||
|
Browser ─┘ └→ agentdb (documents, bills,
|
||||||
|
sync runs, approvals)
|
||||||
|
Firefly API + agentdb → Claude sessions → Notifier → Telegram / ntfy / email
|
||||||
|
```
|
||||||
|
|
||||||
|
Deduplication: use the source's external ID where one exists (SimpleFIN
|
||||||
|
provides one); otherwise a content hash of date + amount + account + normalized
|
||||||
|
description. Firefly's own duplicate detection is the backstop, never the
|
||||||
|
primary mechanism.
|
||||||
|
|
||||||
|
### agentdb schema (v1)
|
||||||
|
|
||||||
|
- `sync_runs` — one row per connector execution: connector, started, finished, status, counts, error.
|
||||||
|
- `documents` — index of filed records: path, sha256, institution, account, doc type, period, source, received date.
|
||||||
|
- `bills` — calendar: payee, expected amount, due date, funding account, autopay flag, status.
|
||||||
|
- `approvals` — bill-pay approvals: bill, amount, requested at, approved/denied at, Telegram message ref, executed at.
|
||||||
|
|
||||||
|
## Safety and error handling
|
||||||
|
|
||||||
|
- **Idempotent imports.** Dedup keys make every retry safe.
|
||||||
|
- **Escalating failure alerts.** Every connector run logs to `sync_runs`; three
|
||||||
|
consecutive failures for a connector escalate to ntfy.
|
||||||
|
- **Secrets stay server-side.** `.env` on the host only, never in the repo.
|
||||||
|
Aggregator access is read-only.
|
||||||
|
- **Approval integrity.** An approval exists only as an `approvals` row created
|
||||||
|
by Jake's own Telegram interaction. The executor re-checks the approval
|
||||||
|
immediately before acting and refuses stale or amount-mismatched approvals.
|
||||||
|
- **Payment receipts.** Every browser-automation payment run captures
|
||||||
|
screenshots and files them as documents.
|
||||||
|
|
||||||
|
## Documents
|
||||||
|
|
||||||
|
Local folder structure (exact layout specified in the Phase 5 spec):
|
||||||
|
`records/<institution>/<account>/<year>/` with normalized filenames
|
||||||
|
(`2026-08-nfcu-checking-statement.pdf`). Every filed document gets a
|
||||||
|
`documents` row. Cloud sync via rclone; Jake picks the target during Phase 5.
|
||||||
|
|
||||||
|
## Testing
|
||||||
|
|
||||||
|
TDD throughout (superpowers test-driven-development skill). Connectors test
|
||||||
|
against recorded fixtures: captured SimpleFIN JSON, real Venmo/Cash App/PayPal
|
||||||
|
CSV exports, sample bill emails. The Firefly client tests against a throwaway
|
||||||
|
local Firefly via a compose test profile. One end-to-end smoke test proves
|
||||||
|
fixture → dedup → Firefly.
|
||||||
|
|
||||||
|
## Deployment
|
||||||
|
|
||||||
|
The stack runs on a remote host behind Traefik; this repo is the config source.
|
||||||
|
Assumed workflow: SSH to the host, `git pull`, `docker compose up -d`. Confirm
|
||||||
|
the exact workflow during Phase 1 planning and record it in the Phase 1 spec.
|
||||||
|
|
||||||
|
## Phases
|
||||||
|
|
||||||
|
Each phase is its own sub-project: spec → plan → implementation. Later phases
|
||||||
|
may not start until the prior phase runs in production.
|
||||||
|
|
||||||
|
1. **Foundation** — Firefly configured for API use (personal access token,
|
||||||
|
asset accounts created); `agent/` skeleton: Dockerfile, config loading,
|
||||||
|
scheduler with a heartbeat job, `agentdb` migrations, compose service,
|
||||||
|
deployed to the host.
|
||||||
|
2. **Ingestion** — SimpleFIN → Firefly sync with dedup; verified coverage for
|
||||||
|
NFCU and TFCU; IMAP receipt connector (promoted from Phase 5) parsing
|
||||||
|
Cash App/Venmo/PayPal transaction emails from the finance inbox; CSV
|
||||||
|
drop-folder connector kept as break-glass fallback.
|
||||||
|
3. **Alerts & monitoring** — notifier with all three sinks; rules engine:
|
||||||
|
low balance, large/unusual transaction, sync failure, negative trend.
|
||||||
|
4. **Analysis & forecasting** — scheduled Claude sessions: weekly spending
|
||||||
|
review, budget tracking, recurring-transaction detection, 30/60/90-day
|
||||||
|
cashflow forecast, monthly digest.
|
||||||
|
5. **Documents** — PDF extraction and filing from the finance inbox (IMAP
|
||||||
|
plumbing already live from Phase 2), browser-automation statement
|
||||||
|
downloads, document index, cloud sync target chosen and wired.
|
||||||
|
6. **Bills** — bill calendar, funding checks, Telegram approve/deny flow,
|
||||||
|
browser automation for approved payments and statement downloads.
|
||||||
|
7. **Later** — business Firefly instance when a venture opens; k8s migration
|
||||||
|
per `k8s/README.md`.
|
||||||
|
|
||||||
|
## Out of scope
|
||||||
|
|
||||||
|
- Investment portfolio management and tax preparation.
|
||||||
|
- Multi-user support; this system serves Jake alone.
|
||||||
|
- Modifying Firefly III itself; the agent is strictly an API consumer.
|
||||||
Reference in New Issue
Block a user