Skip to content
Merged
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
63 changes: 48 additions & 15 deletions src/nanodot/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -181,7 +181,7 @@ def _run_config(args: argparse.Namespace) -> int:
file=sys.stderr,
)
return 1
except (RunnerControlError, OSError) as error:
except (RunnerControlError, OSError, ValueError) as error:
print(f"error: {error}", file=sys.stderr)
return 1
return _run_config_values(args)
Expand All @@ -207,24 +207,48 @@ def _run_config_values(args: argparse.Namespace) -> int:
print("error: a nonempty value is required", file=sys.stderr)
return 1
if is_secret_name(args.name):
store.set(args.name, value)
config.unset(args.name) # remove a legacy plaintext copy after safe save
try:
store.set(args.name, value)
config.unset(args.name) # remove a legacy plaintext copy after safe save
except (ValueError, OSError) as error:
print(f"error: {error}", file=sys.stderr)
return 1
else:
try:
config.set(args.name, value)
except ValueError as error:
except (ValueError, OSError) as error:
print(f"error: {error}", file=sys.stderr)
return 1
elif args.config_command == "unset":
if is_secret_name(args.name):
store.unset(args.name)
config.unset(args.name)
else:
try:
if is_secret_name(args.name):
store.unset(args.name)
config.unset(args.name)
else: # list
rows = [(key, MASK if is_secret_name(key) else str(config.get(key)))
for key in config.keys() if key not in store.names()]
rows += [(name, MASK) for name in store.names()]
except (ValueError, OSError) as error:
print(f"error: {error}", file=sys.stderr)
return 1
else: # list — the diagnosis command must survive a damaged state
try:
keys = config.keys()
secret_names = store.names()
except (ValueError, OSError) as error:
print(f"error: {error}", file=sys.stderr)
return 1
rows: list[tuple[str, str]] = []
for key in keys:
if key in secret_names:
continue
if is_secret_name(key):
# A legacy plaintext copy stored before the secret store
# existed must never be echoed back in the clear.
rows.append((key, MASK))
continue
try:
value = str(config.get(key))
except ValueError as error:
value = f"<invalid: {error}>"
rows.append((key, value))
rows += [(name, MASK) for name in secret_names]
for key, value in sorted(rows):
print(f"{key}={value}")
return 0
Expand All @@ -251,10 +275,11 @@ def _run_watch(args: argparse.Namespace) -> int:

try:
auth_mode = Config().get("github-auth-mode", "token")
except ValueError as error:
has_token = bool(secrets.get("github-token"))
except (ValueError, OSError) as error:
print(f"error: {error}", file=sys.stderr)
return 1
if auth_mode == "token" and not secrets.get("github-token"):
if auth_mode == "token" and not has_token:
print(
"error: no GitHub token configured — run: "
"nanodot config set github-token (hidden prompt)",
Expand Down Expand Up @@ -321,7 +346,12 @@ def _run_watch(args: argparse.Namespace) -> int:
for item in relevant:
print(f" - {item.content}")
if not args.yes:
answer = input("Proceed? [y/N] ").strip().lower()
try:
answer = input("Proceed? [y/N] ").strip().lower()
except EOFError:
# Non-interactive stdin (cron, scripts, closed pipes) must
# decline, not crash — mirroring the secret-prompt path.
answer = "n"
if answer not in ("y", "yes"):
print("cancelled")
return 1
Expand Down Expand Up @@ -420,6 +450,9 @@ def _run_memory(args: argparse.Namespace) -> int:
print("memory is empty")
for item in items:
_print_memory_item(item)
total = memory.count()
if total > len(items):
print(f"... and {total - len(items)} older item(s) not shown")
return 0
except ValueError as error:
print(f"error: {error}", file=sys.stderr)
Expand Down
75 changes: 59 additions & 16 deletions src/nanodot/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,21 @@

Secret-named keys (see redaction.is_secret_name) are routed to the secret
store by the CLI; this file never contains secret material.

Writes serialize through a sidecar flock and replace the file atomically, so
concurrent CLI invocations can neither interleave read-modify-write cycles
(silently reverting a key) nor expose a truncated file to readers.
"""

from __future__ import annotations

import fcntl
import json
import os
import tempfile
from contextlib import contextmanager
from pathlib import Path
from typing import Iterator

from nanodot.paths import data_home

Expand All @@ -27,16 +36,53 @@ def _validated_value(key: str, value: object) -> object:
return value


@contextmanager
def _exclusive_lock(path: Path) -> Iterator[None]:
"""Cross-process mutual exclusion for read-modify-write cycles.

The lock file is never unlinked: an unlinked lock would let two
processes hold two different "locks" on the same path.
"""
with open(path, "a+b") as handle:
fcntl.flock(handle.fileno(), fcntl.LOCK_EX)
try:
yield
finally:
fcntl.flock(handle.fileno(), fcntl.LOCK_UN)


def _atomic_write(path: Path, text: str) -> None:
fd, temporary = tempfile.mkstemp(
prefix=f".{path.name}-", suffix=".tmp", dir=path.parent,
)
try:
with os.fdopen(fd, "w") as fh:
fh.write(text)
fh.flush()
os.fsync(fh.fileno())
os.replace(temporary, path)
finally:
try:
os.unlink(temporary)
except FileNotFoundError:
pass


class Config:
def __init__(self, base: Path | None = None) -> None:
self._path = (base or data_home()) / CONFIG_FILE
self._lock = self._path.parent / f"{self._path.name}.lock"

def get(self, key: str, default: object = None) -> object:
def _read(self) -> dict[str, object]:
if not self._path.exists():
return default
return {}
values = json.loads(self._path.read_text())
if not isinstance(values, dict):
raise ValueError("config.json must contain a JSON object")
return values

def get(self, key: str, default: object = None) -> object:
values = self._read()
if key not in values:
return default
return _validated_value(key, values[key])
Expand All @@ -45,21 +91,18 @@ def set(self, key: str, value: object) -> None:
if key == "mode" and value != "readonly":
raise ValueError("only readonly mode is available; gated/auto modes are not implemented")
value = _validated_value(key, value)
values: dict[str, object] = {}
if self._path.exists():
values = json.loads(self._path.read_text())
if not isinstance(values, dict):
raise ValueError("config.json must contain a JSON object")
values[key] = value
self._path.write_text(json.dumps(values, indent=2))
with _exclusive_lock(self._lock):
values = self._read()
values[key] = value
_atomic_write(self._path, json.dumps(values, indent=2))

def unset(self, key: str) -> None:
if self._path.exists():
values = json.loads(self._path.read_text())
values.pop(key, None)
self._path.write_text(json.dumps(values, indent=2))
with _exclusive_lock(self._lock):
values = self._read()
if key not in values:
return # nothing to remove; never materialize an empty file
values.pop(key)
_atomic_write(self._path, json.dumps(values, indent=2))

def keys(self) -> list[str]:
if not self._path.exists():
return []
return sorted(json.loads(self._path.read_text()))
return sorted(self._read())
2 changes: 1 addition & 1 deletion src/nanodot/core/github_eval.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from nanodot.ports.github import CheckRun, Snapshot

PASSING_CONCLUSIONS = frozenset({"success"})
FAILING_CONCLUSIONS = frozenset({"failure", "timed_out", "cancelled", "action_required", "stale", "error"})
FAILING_CONCLUSIONS = frozenset({"failure", "timed_out", "cancelled", "action_required", "stale", "startup_failure", "error"})
IGNORED_CONCLUSIONS = frozenset({"skipped", "neutral"})


Expand Down
41 changes: 38 additions & 3 deletions src/nanodot/core/memory.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@

import json
import sqlite3
import string
import threading
import time
import uuid
Expand All @@ -49,6 +50,28 @@
DEFAULT_PROPOSAL_EXPIRY_SECONDS = 14 * 24 * 3600 # 14 days


def _require_content(content: str) -> None:
"""Empty items can never match anything and would persist forever."""
if not content or not content.strip():
raise ValueError("memory content must not be empty")


def _keywords(text: str) -> set[str]:
"""One tokenizer for both query and content sides of relevant_to.

Recorded observations carry the PR identity as 'owner/repo#12: ...', so
splitting on '#' and trimming edge punctuation must happen identically
on both sides or the one token guaranteed shared — the target itself —
can never match.
"""
words = set()
for token in text.replace("#", " ").split():
trimmed = token.strip(string.punctuation)
if len(trimmed) > 3:
words.add(trimmed.lower())
return words


@dataclass(frozen=True)
class MemoryItem:
id: str
Expand Down Expand Up @@ -105,6 +128,7 @@ def add_user(
self, content: str, kind: str = KIND_PREFERENCE, at: float | None = None
) -> MemoryItem:
"""Write path 1: the user stated it. Enters confirmed."""
_require_content(content)
return self._insert(
kind=kind,
content=content,
Expand All @@ -123,6 +147,7 @@ def add_observation(
"""Write path 2: an evidenced terminal outcome, auto-recorded by
the runner. Observation, confirmed-by-evidence, provenance links
the evidence."""
_require_content(content)
return self._insert(
kind=KIND_OBSERVATION,
content=content,
Expand All @@ -141,6 +166,7 @@ def propose(
) -> MemoryItem:
"""Write path 3: a proposal (e.g. from a model). Proposed only —
confirmation is a separate, user-driven act."""
_require_content(content)
now = at if at is not None else self._clock.time()
return self._insert(
kind=kind,
Expand Down Expand Up @@ -178,6 +204,7 @@ def confirm(self, item_id: str, at: float | None = None) -> MemoryItem:
return replace(item, status=STATUS_CONFIRMED, expires_at=None)

def edit(self, item_id: str, content: str) -> MemoryItem:
_require_content(content)
self._require(item_id)
with self._lock:
self._conn.execute(
Expand Down Expand Up @@ -251,18 +278,26 @@ def list(
).fetchall()
return [self._row_to_item(row) for row in rows]

def count(self, status: str | None = None) -> int:
self.sweep_expired()
clause, params = ("WHERE status=?", [status]) if status else ("", [])
with self._lock:
row = self._conn.execute(
f"SELECT COUNT(*) FROM memory {clause}", params
).fetchone()
return int(row[0])

def context_for_prompts(self) -> list[str]:
"""Confirmed content only — what the local read path may consult."""
return [item.content for item in self.list(status=STATUS_CONFIRMED)]

def relevant_to(self, text: str) -> list[MemoryItem]:
"""Confirmed items whose content shares words with the text — the
local 'remembered context' surfaced at task creation."""
words = {w.lower() for w in text.replace("#", " ").split() if len(w) > 3}
words = _keywords(text)
out = []
for item in self.list(status=STATUS_CONFIRMED):
item_words = {w.lower() for w in item.content.split() if len(w) > 3}
if item_words & words:
if _keywords(item.content) & words:
out.append(item)
return out

Expand Down
4 changes: 4 additions & 0 deletions src/nanodot/core/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,10 @@ def validate(self, task: Task) -> None:

def create(self, task: Task) -> Task:
self.validate(task)
if task.next_check_at is None:
# SQL `NULL <= now` is never true: persisted without a schedule,
# an ACTIVE task would silently never be fetched.
task.next_check_at = self._time.time()
try:
with self._lock, self._conn:
self._conn.execute(
Expand Down
11 changes: 10 additions & 1 deletion src/nanodot/native/daemon.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,16 @@ def tick(self, stop: threading.Event | None = None) -> int:
"""
now = self._clock.time()
attempted = 0
for task in self._store.list_schedulable(now):
try:
schedulable = self._store.list_schedulable(now)
except Exception:
# Store iteration runs outside the per-task guard below. A
# transient failure (secret rotation race, momentary lock) must
# cost this pass, not the daemon; the next tick retries. Never
# log exception text: it can contain credentials or private data.
logger.warning("task listing failed; will retry next pass")
return 0
for task in schedulable:
if stop is not None and stop.is_set():
break
attempted += 1
Expand Down
Loading
Loading