Update projektu
This commit is contained in:
@@ -1,316 +0,0 @@
|
||||
#!/usr/bin/env -S uv run --script
|
||||
# /// script
|
||||
# dependencies = []
|
||||
# ///
|
||||
|
||||
"""
|
||||
note.py — backend for /note skill.
|
||||
SQLite-backed note store with tags, soft-delete, and operation log.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import re
|
||||
import sqlite3
|
||||
import sys
|
||||
from collections.abc import Generator
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
DB_PATH = Path(__file__).resolve().parent.parent.parent.parent / "db" / "note.sqlite"
|
||||
LOG_PATH = Path(__file__).resolve().parent.parent.parent.parent / "log" / "note.log"
|
||||
|
||||
_TAG_RE = re.compile(r"^[a-z][a-z0-9-]*$")
|
||||
# A URL together with an immediately preceding "Label:" token, if any.
|
||||
# The leading separator class swallows the connector that introduced the URL
|
||||
# (em-dash, comma, etc.) so it does not dangle once the URL moves to its own line.
|
||||
_LABELED_URL_RE = re.compile(r"[\s,;—–-]*([^\s,]+:\s*)?(https?://[^\s,]+)")
|
||||
|
||||
SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS notes (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
content TEXT NOT NULL,
|
||||
tags TEXT NOT NULL DEFAULT '[]',
|
||||
created_at TEXT NOT NULL,
|
||||
deleted_at TEXT
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS tags (
|
||||
name TEXT PRIMARY KEY,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
"""
|
||||
|
||||
|
||||
def _init_db(conn: sqlite3.Connection) -> None:
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.executescript(SCHEMA)
|
||||
_migrate(conn)
|
||||
|
||||
|
||||
def _migrate(conn: sqlite3.Connection) -> None:
|
||||
cols = {row[1] for row in conn.execute("PRAGMA table_info(notes)")}
|
||||
if "tags" not in cols:
|
||||
conn.execute("ALTER TABLE notes ADD COLUMN tags TEXT NOT NULL DEFAULT '[]'")
|
||||
if "deleted_at" not in cols:
|
||||
conn.execute("ALTER TABLE notes ADD COLUMN deleted_at TEXT")
|
||||
conn.commit()
|
||||
_backfill_tags(conn)
|
||||
|
||||
|
||||
def _backfill_tags(conn: sqlite3.Connection) -> None:
|
||||
"""On first introduction of the registry, seed it from tags already used in notes."""
|
||||
existing = {row[0] for row in conn.execute("SELECT name FROM tags")}
|
||||
if existing:
|
||||
return
|
||||
used = {row[0] for row in conn.execute("SELECT DISTINCT value FROM notes, json_each(notes.tags)")}
|
||||
if not used:
|
||||
return
|
||||
now = datetime.now(timezone.utc).isoformat()
|
||||
conn.executemany(
|
||||
"INSERT OR IGNORE INTO tags(name, created_at) VALUES(?, ?)",
|
||||
[(tag, now) for tag in sorted(used)],
|
||||
)
|
||||
conn.commit()
|
||||
|
||||
|
||||
@contextmanager
|
||||
def _connect() -> Generator[sqlite3.Connection, None, None]:
|
||||
DB_PATH.parent.mkdir(parents=True, exist_ok=True)
|
||||
conn = sqlite3.connect(DB_PATH)
|
||||
conn.row_factory = sqlite3.Row
|
||||
_init_db(conn)
|
||||
try:
|
||||
yield conn
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _validate_tags(tags: list[str]) -> None:
|
||||
for tag in tags:
|
||||
if not _TAG_RE.match(tag):
|
||||
raise ValueError(
|
||||
f"Invalid tag '{tag}' — use lowercase letters, digits, hyphens only (e.g. cli, soft-delete)"
|
||||
)
|
||||
|
||||
|
||||
def _tags_display(tags_json: str) -> str:
|
||||
tags = json.loads(tags_json)
|
||||
if not tags:
|
||||
return ""
|
||||
return " [" + " ".join(f"#{t}" for t in tags) + "]"
|
||||
|
||||
|
||||
def _urls_on_own_lines(text: str) -> str:
|
||||
"""Lay out each URL (and its inline "Label:", if any) on its own bullet line.
|
||||
|
||||
The chat UI merges two adjacent links into one block and hides the second,
|
||||
which also overlays the list number. Putting each URL on its own line keeps
|
||||
them separate and the number visible. URLs stay bare so they autolink.
|
||||
"""
|
||||
if not _LABELED_URL_RE.search(text):
|
||||
return text
|
||||
|
||||
def repl(match: re.Match[str]) -> str:
|
||||
label = match.group(1) or ""
|
||||
return f"\n - {label}{match.group(2)}"
|
||||
|
||||
return _LABELED_URL_RE.sub(repl, text)
|
||||
|
||||
|
||||
def _log(op: str, detail: str) -> None:
|
||||
LOG_PATH.parent.mkdir(parents=True, exist_ok=True)
|
||||
ts = datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]
|
||||
with LOG_PATH.open("a") as f:
|
||||
f.write(f"{ts} {op} {detail}\n")
|
||||
|
||||
|
||||
def _active_ids(conn: sqlite3.Connection) -> list[int]:
|
||||
rows = conn.execute(
|
||||
"SELECT id FROM notes WHERE deleted_at IS NULL ORDER BY created_at DESC"
|
||||
).fetchall()
|
||||
return [row["id"] for row in rows]
|
||||
|
||||
|
||||
def cmd_add(args: argparse.Namespace) -> int:
|
||||
tags: list[str] = args.tags or []
|
||||
try:
|
||||
_validate_tags(tags)
|
||||
except ValueError as exc:
|
||||
print(str(exc), file=sys.stderr)
|
||||
return 1
|
||||
content = args.text.strip()
|
||||
tags_json = json.dumps(tags)
|
||||
created_at = datetime.now(timezone.utc).isoformat()
|
||||
with _connect() as conn:
|
||||
known = {row[0] for row in conn.execute("SELECT name FROM tags")}
|
||||
unknown = [tag for tag in tags if tag not in known]
|
||||
if unknown:
|
||||
print(f"Unknown tag(s): {', '.join(unknown)}", file=sys.stderr)
|
||||
return 2
|
||||
cur = conn.execute(
|
||||
"INSERT INTO notes(content, tags, created_at) VALUES(?, ?, ?)",
|
||||
(content, tags_json, created_at),
|
||||
)
|
||||
conn.commit()
|
||||
nid = cur.lastrowid
|
||||
tags_log = ",".join(tags)
|
||||
_log("ADD", f"id={nid} tags=[{tags_log}] {content}")
|
||||
print(f"Noted [#1]: {content}{_tags_display(tags_json)}")
|
||||
return 0
|
||||
|
||||
|
||||
def cmd_list(args: argparse.Namespace) -> int:
|
||||
with _connect() as conn:
|
||||
id_to_display = {nid: i + 1 for i, nid in enumerate(_active_ids(conn))}
|
||||
if args.tag:
|
||||
placeholders = ",".join("?" * len(args.tag))
|
||||
rows = conn.execute(
|
||||
f"""
|
||||
SELECT id, content, tags FROM notes
|
||||
WHERE deleted_at IS NULL
|
||||
AND (
|
||||
SELECT count(*) FROM json_each(notes.tags)
|
||||
WHERE value IN ({placeholders})
|
||||
) > 0
|
||||
ORDER BY created_at DESC
|
||||
LIMIT ? OFFSET ?
|
||||
""",
|
||||
(*args.tag, args.limit, args.offset),
|
||||
).fetchall()
|
||||
else:
|
||||
rows = conn.execute(
|
||||
"SELECT id, content, tags FROM notes"
|
||||
" WHERE deleted_at IS NULL"
|
||||
" ORDER BY created_at DESC LIMIT ? OFFSET ?",
|
||||
(args.limit, args.offset),
|
||||
).fetchall()
|
||||
tag_filter = ",".join(args.tag) if args.tag else "None"
|
||||
_log("LIST", f"tag={tag_filter} returned={len(rows)}")
|
||||
if not rows:
|
||||
print("No notes.")
|
||||
return 0
|
||||
for row in rows:
|
||||
head, sep, rest = _urls_on_own_lines(row["content"]).partition("\n")
|
||||
print(f"{id_to_display[row['id']]}. {head}{_tags_display(row['tags'])}{sep}{rest}")
|
||||
return 0
|
||||
|
||||
|
||||
def cmd_delete(args: argparse.Namespace) -> int:
|
||||
display_id: int = args.id
|
||||
deleted_at = datetime.now(timezone.utc).isoformat()
|
||||
with _connect() as conn:
|
||||
ids = _active_ids(conn)
|
||||
idx = display_id - 1
|
||||
if idx < 0 or idx >= len(ids):
|
||||
print(f"No active note with display id={display_id}.")
|
||||
return 1
|
||||
nid = ids[idx]
|
||||
row = conn.execute(
|
||||
"SELECT id, content, tags FROM notes WHERE id = ?", (nid,)
|
||||
).fetchone()
|
||||
conn.execute("UPDATE notes SET deleted_at = ? WHERE id = ?", (deleted_at, nid))
|
||||
conn.commit()
|
||||
tags_log = ",".join(json.loads(row["tags"]))
|
||||
_log("DELETE", f"display_id={display_id} id={nid} tags=[{tags_log}] content={row['content']!r}")
|
||||
print(f"Deleted: {row['content']}")
|
||||
return 0
|
||||
|
||||
|
||||
def cmd_show(args: argparse.Namespace) -> int:
|
||||
display_id: int = args.id
|
||||
with _connect() as conn:
|
||||
ids = _active_ids(conn)
|
||||
idx = display_id - 1
|
||||
if idx < 0 or idx >= len(ids):
|
||||
print(f"No active note with display id={display_id}.")
|
||||
return 1
|
||||
nid = ids[idx]
|
||||
row = conn.execute(
|
||||
"SELECT id, content, tags, created_at FROM notes WHERE id = ?", (nid,)
|
||||
).fetchone()
|
||||
_log("SHOW", f"display_id={display_id} id={nid}")
|
||||
tags = json.loads(row["tags"])
|
||||
tags_line = " ".join(f"#{t}" for t in tags) if tags else "(none)"
|
||||
print(f"Note [#{display_id}] (id={row['id']})")
|
||||
print(f"created: {row['created_at']}")
|
||||
print(f"tags: {tags_line}")
|
||||
print(f"content: {row['content']}")
|
||||
return 0
|
||||
|
||||
|
||||
def cmd_tag_add(args: argparse.Namespace) -> int:
|
||||
name = args.name.strip()
|
||||
try:
|
||||
_validate_tags([name])
|
||||
except ValueError as exc:
|
||||
print(str(exc), file=sys.stderr)
|
||||
return 1
|
||||
with _connect() as conn:
|
||||
exists = conn.execute("SELECT 1 FROM tags WHERE name = ?", (name,)).fetchone()
|
||||
if exists:
|
||||
print(f"Tag '#{name}' already exists.")
|
||||
return 0
|
||||
created_at = datetime.now(timezone.utc).isoformat()
|
||||
conn.execute("INSERT INTO tags(name, created_at) VALUES(?, ?)", (name, created_at))
|
||||
conn.commit()
|
||||
_log("TAG-ADD", f"name={name}")
|
||||
print(f"Tag created: #{name}")
|
||||
return 0
|
||||
|
||||
|
||||
def cmd_tag_list(args: argparse.Namespace) -> int:
|
||||
with _connect() as conn:
|
||||
rows = conn.execute("SELECT name FROM tags ORDER BY name").fetchall()
|
||||
_log("TAG-LIST", f"returned={len(rows)}")
|
||||
if not rows:
|
||||
print("No tags.")
|
||||
return 0
|
||||
for row in rows:
|
||||
print(f"#{row['name']}")
|
||||
return 0
|
||||
|
||||
|
||||
def _main() -> int:
|
||||
parser = argparse.ArgumentParser(description="Note store")
|
||||
sub = parser.add_subparsers(dest="cmd", required=True)
|
||||
|
||||
p_add = sub.add_parser("add", help="Add a note")
|
||||
p_add.add_argument("text", help="Note content")
|
||||
p_add.add_argument("--tags", nargs="+", metavar="TAG", default=[], help="Tags (lowercase, hyphens allowed)")
|
||||
|
||||
p_list = sub.add_parser("list", help="List active notes")
|
||||
p_list.add_argument("--limit", type=int, default=50)
|
||||
p_list.add_argument("--offset", type=int, default=0)
|
||||
p_list.add_argument("--tag", nargs="+", metavar="TAG", help="Filter by tag (OR logic)")
|
||||
|
||||
p_show = sub.add_parser("show", help="Show one note in full by display ID")
|
||||
p_show.add_argument("id", type=int, help="Display ID")
|
||||
|
||||
p_del = sub.add_parser("delete", help="Soft-delete a note by ID")
|
||||
p_del.add_argument("id", type=int, help="Note ID")
|
||||
|
||||
p_tag_add = sub.add_parser("tag-add", help="Register a tag")
|
||||
p_tag_add.add_argument("name", help="Tag name (lowercase, hyphens allowed)")
|
||||
|
||||
sub.add_parser("tag-list", help="List registered tags")
|
||||
|
||||
args = parser.parse_args()
|
||||
|
||||
if args.cmd == "add":
|
||||
return cmd_add(args)
|
||||
if args.cmd == "list":
|
||||
return cmd_list(args)
|
||||
if args.cmd == "show":
|
||||
return cmd_show(args)
|
||||
if args.cmd == "delete":
|
||||
return cmd_delete(args)
|
||||
if args.cmd == "tag-add":
|
||||
return cmd_tag_add(args)
|
||||
if args.cmd == "tag-list":
|
||||
return cmd_tag_list(args)
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(_main())
|
||||
111
skills/note/scripts/note_capture.py
Normal file
111
skills/note/scripts/note_capture.py
Normal file
@@ -0,0 +1,111 @@
|
||||
#!/usr/bin/env -S uv run --script
|
||||
# /// script
|
||||
# requires-python = ">=3.11"
|
||||
# dependencies = []
|
||||
# ///
|
||||
"""note_capture.py — dumb, instant capture for the /note skill.
|
||||
|
||||
Writes the raw input verbatim into notes/inbox/ (atomic tmp -> os.replace) plus one
|
||||
audit line to log/note.log, then prints a one-line confirmation. No reformulation,
|
||||
no reading of the knowledge doc, no compile — that is the compile step's job
|
||||
(inline in immediate mode, or the cron drain in `cron` mode).
|
||||
|
||||
Used identically by both modes; the only difference is what the agent does *after*
|
||||
calling this (immediate: run the compile workflow inline; cron: stop).
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import re
|
||||
import sys
|
||||
import unicodedata
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
# workspace/skills/note/scripts/note_capture.py -> parents[3] = workspace root.
|
||||
WORKSPACE = Path(__file__).resolve().parents[3]
|
||||
INBOX = WORKSPACE / "notes" / "inbox"
|
||||
LOG = WORKSPACE / "log" / "note.log"
|
||||
|
||||
_SLUG_STRIP_RE = re.compile(r"[^a-z0-9]+")
|
||||
_URL_RE = re.compile(r"https?://([^/\s]+)")
|
||||
MAX_SLUG_WORDS = 4
|
||||
SLUG_MAX_LEN = 40
|
||||
|
||||
|
||||
def _ascii_fold(text: str) -> str:
|
||||
"""Drop diacritics so Czech words survive slugging (mazání -> mazani)."""
|
||||
return unicodedata.normalize("NFKD", text).encode("ascii", "ignore").decode("ascii")
|
||||
|
||||
|
||||
def _slugify(text: str) -> str:
|
||||
"""Short kebab slug from the first words of the input (domain for a bare URL)."""
|
||||
first_line = next((line for line in text.splitlines() if line.strip()), "").strip()
|
||||
url_match = _URL_RE.match(first_line)
|
||||
if url_match:
|
||||
host = url_match.group(1).removeprefix("www.")
|
||||
slug = _SLUG_STRIP_RE.sub("-", _ascii_fold(host).lower()).strip("-")
|
||||
return slug or "note"
|
||||
words = first_line.split()[:MAX_SLUG_WORDS]
|
||||
slug = _SLUG_STRIP_RE.sub("-", _ascii_fold(" ".join(words)).lower()).strip("-")
|
||||
return slug[:SLUG_MAX_LEN].strip("-") or "note"
|
||||
|
||||
|
||||
def _build_content(
|
||||
captured_at: str, channel: str | None, chat_id: str | None, body: str
|
||||
) -> str:
|
||||
lines = [f"captured_at: {captured_at}"]
|
||||
if channel:
|
||||
lines.append(f"channel: {channel}")
|
||||
if chat_id:
|
||||
lines.append(f'chat_id: "{chat_id}"')
|
||||
frontmatter = "\n".join(lines)
|
||||
return f"---\n{frontmatter}\n---\n\n{body.strip()}\n"
|
||||
|
||||
|
||||
def _append_log(filename: str, body: str) -> None:
|
||||
LOG.parent.mkdir(parents=True, exist_ok=True)
|
||||
stamp = datetime.now().astimezone().isoformat(timespec="seconds")
|
||||
summary = " ".join(body.split())[:80]
|
||||
with LOG.open("a", encoding="utf-8") as handle:
|
||||
handle.write(f"{stamp} CAPTURE {filename} :: {summary}\n")
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(description="Capture a raw note into notes/inbox/")
|
||||
parser.add_argument(
|
||||
"--text", default=None, help="Raw input; if omitted, read from stdin"
|
||||
)
|
||||
parser.add_argument(
|
||||
"--channel", default=None, help="Origin channel (telegram/websocket/cli)"
|
||||
)
|
||||
parser.add_argument(
|
||||
"--chat-id", default=None, dest="chat_id", help="Origin chat id"
|
||||
)
|
||||
args = parser.parse_args()
|
||||
|
||||
body = args.text if args.text is not None else sys.stdin.read()
|
||||
body = body.strip()
|
||||
if not body:
|
||||
print("Nothing to capture (empty input).", file=sys.stderr)
|
||||
return 1
|
||||
|
||||
now = datetime.now().astimezone()
|
||||
timestamp = now.strftime("%Y-%m-%d_%H_%M_%S_%f")
|
||||
filename = f"{timestamp}-{_slugify(body)}.md"
|
||||
content = _build_content(
|
||||
now.isoformat(timespec="seconds"), args.channel, args.chat_id, body
|
||||
)
|
||||
|
||||
INBOX.mkdir(parents=True, exist_ok=True)
|
||||
tmp_path = INBOX / f".{filename}.tmp"
|
||||
final_path = INBOX / filename
|
||||
tmp_path.write_text(content, encoding="utf-8")
|
||||
tmp_path.replace(final_path)
|
||||
|
||||
_append_log(filename, body)
|
||||
print(f"captured: {filename}")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
216
skills/note/scripts/note_compile.py
Normal file
216
skills/note/scripts/note_compile.py
Normal file
@@ -0,0 +1,216 @@
|
||||
#!/usr/bin/env -S uv run --script
|
||||
# /// script
|
||||
# requires-python = ">=3.11"
|
||||
# dependencies = ["nanobot-ai"]
|
||||
# ///
|
||||
"""note_compile.py — drain notes/inbox/ into the structured doc notes/notes.md.
|
||||
|
||||
Thin launcher run by the nanobot user crontab every minute. All the intelligence
|
||||
lives in DRAIN_GOAL + the note skill's compile workflow; this script only decides
|
||||
*when* to run and guards against concurrent runs.
|
||||
|
||||
Flow:
|
||||
1. Cheap fs pre-check (no LLM): are there pending files in notes/inbox/? None ->
|
||||
exit 0 without importing nanobot (per-minute polling stays nearly free).
|
||||
2. Lockfile (notes/.compile.lock, PID + start-timestamp): another compile running?
|
||||
-> exit 0. Stale lock (dead PID / older than STALE_SECONDS) is reclaimed.
|
||||
3. Otherwise Nanobot.from_config().run(<drain goal>) — drains ALL pending in one
|
||||
batch. process_direct has NO cron preamble (unlike cron/jobs.json agent jobs).
|
||||
4. Quietly append to log/note_compile_cron.log; no Telegram.
|
||||
|
||||
The immediate `/note` mode runs the SAME compile workflow inline and takes the SAME
|
||||
lock, so an inline merge and a background drain cannot corrupt notes.md at once.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import traceback
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
# workspace/skills/note/scripts/note_compile.py -> parents[3] = workspace root.
|
||||
WORKSPACE = Path(__file__).resolve().parents[3]
|
||||
NOTES = WORKSPACE / "notes"
|
||||
INBOX = NOTES / "inbox"
|
||||
LOCK = NOTES / ".compile.lock"
|
||||
LOG = WORKSPACE / "log" / "note_compile_cron.log"
|
||||
|
||||
TIMEOUT_SECONDS = 15 * 60
|
||||
STALE_SECONDS = 30 * 60
|
||||
|
||||
DRAIN_GOAL = (
|
||||
"Pomocí skillu note (Compile/drain) zpracuj VŠECHNY čekající soubory v `notes/inbox/` "
|
||||
"(regulérní soubory přímo v `notes/inbox/`, mimo skryté). Pro každý postupuj podle "
|
||||
"*Compile workflow* v note SKILL.md: přeformuluj na terse fakt(a) (zachovej jazyk vstupu, "
|
||||
"jeden koncept per záznam, zahoď filler); z těla vytáhni VŠECHNY URL (0..N) a každou stáhni "
|
||||
"přes `web` tool — když je za paywallem / login-wallem / neúplná, NEfabrikuj shrnutí, zapiš "
|
||||
"jen URL + titulek + značku `⚠ paywall/neúplné`. Zařaď obsah pod správnou tematickou sekci "
|
||||
"v `notes/notes.md` (novou sekci ## založ, když chybí; existující sekci uprav chirurgicky, "
|
||||
"nepřepisuj celý dokument). Po úspěšném zařazení přesuň zdrojový soubor do `notes/done/`; "
|
||||
"když z něj nešlo nic použitelného získat (vše za paywallem / nečitelné / nejednoznačné), "
|
||||
"přesuň ho do `notes/hard/`. Přesouvej HNED po každém souboru, ať ho příští cron tik "
|
||||
"nezpracovává znovu. Běžíš v izolované session na pozadí, bez interakce s uživatelem."
|
||||
)
|
||||
|
||||
|
||||
def log(message: str) -> None:
|
||||
LOG.parent.mkdir(parents=True, exist_ok=True)
|
||||
stamp = datetime.now().astimezone().isoformat(timespec="seconds")
|
||||
with LOG.open("a", encoding="utf-8") as handle:
|
||||
handle.write(f"{stamp} {message}\n")
|
||||
|
||||
|
||||
def pending_sources() -> list[Path]:
|
||||
"""Regular files directly in notes/inbox/ (hidden files excluded; done/ and hard/ are siblings)."""
|
||||
if not INBOX.exists():
|
||||
return []
|
||||
return [
|
||||
p for p in sorted(INBOX.iterdir()) if p.is_file() and not p.name.startswith(".")
|
||||
]
|
||||
|
||||
|
||||
def _pid_alive(pid: int) -> bool:
|
||||
try:
|
||||
os.kill(pid, 0)
|
||||
except ProcessLookupError:
|
||||
return False
|
||||
except PermissionError:
|
||||
return True
|
||||
return True
|
||||
|
||||
|
||||
def _lock_is_stale() -> bool:
|
||||
"""A lock is dead if unreadable, its PID is gone, or it is older than STALE_SECONDS."""
|
||||
try:
|
||||
data = json.loads(LOCK.read_text())
|
||||
pid = int(data["pid"])
|
||||
started = datetime.fromisoformat(data["started"])
|
||||
except (OSError, ValueError, KeyError):
|
||||
return True
|
||||
if not _pid_alive(pid):
|
||||
return True
|
||||
age = (datetime.now().astimezone() - started).total_seconds()
|
||||
return age > STALE_SECONDS
|
||||
|
||||
|
||||
def acquire_lock() -> bool:
|
||||
"""Atomically create the lock. Return False when a live compile already runs."""
|
||||
for _ in range(2):
|
||||
try:
|
||||
fd = os.open(LOCK, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o644)
|
||||
except FileExistsError:
|
||||
if not _lock_is_stale():
|
||||
return False
|
||||
log("stale lock, reclaiming")
|
||||
LOCK.unlink(missing_ok=True)
|
||||
continue
|
||||
payload = {
|
||||
"pid": os.getpid(),
|
||||
"started": datetime.now().astimezone().isoformat(),
|
||||
}
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as handle:
|
||||
json.dump(payload, handle)
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
async def run_compile(goal: str) -> str:
|
||||
# Heavy import deferred: the per-minute pre-check (no pending work) must not pay the
|
||||
# nanobot import cost — only an actual compile run needs it.
|
||||
from nanobot import Nanobot
|
||||
|
||||
bot = Nanobot.from_config()
|
||||
result = await bot.run(goal, session_key="note-compile")
|
||||
return result.content or ""
|
||||
|
||||
|
||||
def commit_notes(count: int) -> None:
|
||||
"""Stage and commit only notes/ after a successful drain.
|
||||
|
||||
The Dream processor owns the rest of the workspace, so we never `git add -A`.
|
||||
A no-op when notes/ has no changes. .compile.lock is gitignored, so `git add notes/`
|
||||
(run while the lock is still held) does not stage it. Commit failure is logged, not
|
||||
raised — the drain itself already succeeded and must not be reported as failed.
|
||||
"""
|
||||
try:
|
||||
status = subprocess.run(
|
||||
["git", "-C", str(WORKSPACE), "status", "--porcelain", "notes/"],
|
||||
check=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
if not status.stdout.strip():
|
||||
return
|
||||
subprocess.run(
|
||||
["git", "-C", str(WORKSPACE), "add", "notes/"],
|
||||
check=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
subprocess.run(
|
||||
[
|
||||
"git",
|
||||
"-C",
|
||||
str(WORKSPACE),
|
||||
"commit",
|
||||
"-m",
|
||||
f"note: cron drain ({count} captures)",
|
||||
],
|
||||
check=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
)
|
||||
log(f"COMMIT notes/ ({count} captures)")
|
||||
except (OSError, subprocess.CalledProcessError) as error:
|
||||
log(f"WARN commit failed: {error}")
|
||||
|
||||
|
||||
def main() -> int:
|
||||
dry_run = "--dry-run" in sys.argv[1:]
|
||||
|
||||
pending = pending_sources()
|
||||
if not pending:
|
||||
return 0
|
||||
|
||||
if not acquire_lock():
|
||||
log(f"SKIP compile already running ({len(pending)} pending)")
|
||||
return 0
|
||||
|
||||
if dry_run:
|
||||
names = ", ".join(p.name for p in pending)
|
||||
log(f"DRY-RUN would compile {len(pending)} pending: {names}")
|
||||
LOCK.unlink(missing_ok=True)
|
||||
return 0
|
||||
|
||||
started = datetime.now().astimezone()
|
||||
log(f"START compile {len(pending)} pending: {', '.join(p.name for p in pending)}")
|
||||
try:
|
||||
result_text = asyncio.run(
|
||||
asyncio.wait_for(run_compile(DRAIN_GOAL), timeout=TIMEOUT_SECONDS)
|
||||
)
|
||||
summary = (
|
||||
result_text.strip().splitlines()[0][:200]
|
||||
if result_text.strip()
|
||||
else "(prázdný výstup)"
|
||||
)
|
||||
duration = int((datetime.now().astimezone() - started).total_seconds())
|
||||
log(
|
||||
f"END compile duration={duration}s remaining={len(pending_sources())} :: {summary}"
|
||||
)
|
||||
commit_notes(len(pending))
|
||||
return 0
|
||||
except asyncio.TimeoutError:
|
||||
log(f"TIMEOUT compile po {TIMEOUT_SECONDS // 60} min")
|
||||
return 1
|
||||
except Exception as error:
|
||||
log(f"EXCEPTION compile: {error}\n{traceback.format_exc()}")
|
||||
return 1
|
||||
finally:
|
||||
LOCK.unlink(missing_ok=True)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
Reference in New Issue
Block a user