diff --git a/.gitignore b/.gitignore index 354fecb..5b97cac 100644 --- a/.gitignore +++ b/.gitignore @@ -36,3 +36,7 @@ packages/core/bench/agent-recall-results.json # MCP publisher credentials (some publisher versions store login in CWD) .mcpregistry_* + +# Local release preparation tools and install smoke checks +.release-tools/ +.release-smoke/ diff --git a/CHANGELOG.md b/CHANGELOG.md index c633c5d..1d58b84 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,21 @@ # Changelog +## 0.16.0 - Complete source reads, safe retries and fact freshness + +Unreleased. + +- **Scoped native search.** GitHub and Notion search follow result pages and retain configured repository/database scope. GitHub's 1000-match cap, server incompleteness, and missing/repeated Notion cursors become explicit search warnings. +- **HTTP deadlines and cancellation.** Shared HTTP requests use a 30-second deadline through retries and response-body consumption (`ATS_HTTP_TIMEOUT_MS`). Caller cancellation stops retry waits and prevents another attempt. +- **Per-project corpus failures.** GitHub and Notion keep healthy project records when another project fails, with source warnings that propagate through retrieval. Partial corpora remain uncached. +- **Reviewed fact freshness.** `ats kg stale --days N` lists active facts due for review. `ats kg confirm FACT --source REF` stages new evidence for approval; ratification appends a confirmation without changing original validity or provenance. +- **Safe mutation retries.** Reads retain transient retry behavior. Mutating requests retry only explicit rate-limit rejections; ambiguous network and gateway failures return without replay. Adapter read-only POST requests can opt in to safe retries. Long server reset times return the rejection instead of retrying too early. +- **Fact ratification integrity.** Fact writes verify the approved payload, serialize store checks with the append, recheck conflicts, and deduplicate repeated ratification by proposal id. CLI ratification claims the approved review item before writing. +- **Bound create keys.** Idempotency keys bind the effective request to its source configuration and acquire a durable claim before writing, including reviewed creates. Changed bindings, concurrent attempts and uncertain outcomes fail closed. +- **Bound batch resumes.** Journals bind each item to its payload and source, record an applying claim before execution, and skip only matching completed operations. Changed, interrupted, failed or legacy unbound entries require inspection. Dry runs preserve journal bytes. +- **Source fact dates.** `kg propose --valid-at ISO --learned-at ISO` records when a fact happened and when it was learned separately from ratification. JSONL accepts `validAt`/`learnedAt`; timestamps normalize to UTC and future/planned dates are refused. Cypher and Graphiti exports retain both times. + + ## 0.15.0 - Reviewed writes and CLI state integrity Released 2026-10-02. diff --git a/README.md b/README.md index eac51e3..0299bed 100644 --- a/README.md +++ b/README.md @@ -97,6 +97,8 @@ ATS is a good fit when operational context already lives in task systems or conn ATS-managed execution metadata can be encoded in the task body, with typed links under `## Related` and consulted sources under `## References`. Managed helpers are designed to preserve human-authored rows and links; `update --content` replaces the complete body, so callers that add to a body use `--append` / `--prepend`, present the `contentHash` they read with `--if-match`, and verify the result. [`npm run prove:intent`](examples/intent-layer/) runs a deterministic synthetic proof of the execution-context path. +For reliable automation and freshness workflows, see [CLI reliability](docs/cli-reliability.md): scoped pagination, bounded requests, payload-bound retry/resume keys, and reviewed fact confirmations. + ## Minimal adapter-neutral workflow 1. **Select and verify an adapter:** `ats config use `, authenticate as its README describes, then run `ats doctor`. diff --git a/docs/cli-reliability.md b/docs/cli-reliability.md new file mode 100644 index 0000000..649d2ba --- /dev/null +++ b/docs/cli-reliability.md @@ -0,0 +1,43 @@ +# CLI reliability + +## Complete reads + +GitHub and Notion native search follow result pages and retain the adapter's configured repository or database allow-list, including when the query itself names another repository. GitHub limits search to 1000 matches and can report incomplete results. Notion search stops with a warning on missing/repeated cursors or after 100 pages. Narrow the query when a warning reports a bound. + +A failed GitHub repository or Notion database leaves healthy corpus records available. The result carries per-source warnings; Core does not cache partial corpora. `ats find QUERY --require-complete --json` preserves its diagnostic output and returns exit 2 for degraded reads. + +## Bounded requests and safe retries + +`ATS_HTTP_TIMEOUT_MS` is a positive millisecond deadline, default 30000, spanning request attempts, retry waits, and response-body reads. The shared HTTP helper accepts `signal` and forwards cancellation to fetch. A cancelled wait cannot start another request. Large server reset times return the rate-limit response instead of retrying before the requested reset. + +Reads retry transient network/gateway failures. Mutations retry explicit rate-limit rejection only. An ambiguous gateway/network failure might follow a successful remote write, so inspect the backend before repeating it. Adapter authors can mark a read-only POST with `retrySafe: true`; Notion search and database queries use this contract. + +## Create keys and batch journals + +```sh +ats create writing "Release checklist" --idempotency-key release-checklist --json +ats batch changes.jsonl --journal changes.journal.jsonl --json +``` + +A create key binds the effective normalized payload and source scope. A matching completed request replays its recorded task or review item. A changed payload or source returns exit 3. The claim is durable before a create, including applying a reviewed create; concurrent callers cannot acquire it again. In-flight/uncertain claims never expire into automatic duplicates. + +A journal binds each item id to its full operation and source scope. An applying claim precedes execution. A matching applied/staged operation is skipped on resume; altered, failed, or unfinished operations fail with per-item diagnostics and batch exit 5. Dry-run executes validation/planning without writing or claiming the journal. Journaled claims also exclude concurrent processes. + +Claims are intentionally conservative. A crash after a backend write and before local recording can leave an uncertain outcome. Inspect the source and local receipt before staging fresh work. Older keys/journals without payload bindings cannot safely resume: inspect them, then use a fresh key or journal. Source scope includes configuration and working directory; changing either can require a fresh binding. Claims do not provide a distributed transaction with the backend. + +## Fact dates and freshness + +```sh +ats kg propose "Acme" "uses" "Invoice service" --source "source-record" \ + --valid-at 2026-01-01T00:00:00Z --learned-at 2026-02-01T00:00:00Z --json +ats kg stale --days 60 --domain sales --json +ats kg confirm FACT_ID --source "verified source record" --json +ats review approve REVIEW_ID +ats kg ratify REVIEW_ID --json +``` + +Explicit source times must be ISO timestamps with a timezone, normalize to UTC and cannot be future/planned dates. JSONL proposals accept `validAt` and `learnedAt`. Ratified facts carry `tValid`, `tLearned` and `provenance.ratifiedAt`; omitted valid time retains ratification-time behavior, while omitted learned time uses staging time. `--as-of` continues to query validity intervals; learned time discloses when historical evidence entered the workflow. Exports retain both times. + +Freshness age uses the last reviewed confirmation, then learned/ratified time, then legacy validity time. Age is a review signal and never automatically closes a fact. Confirmation requires source evidence, approval and ratification, then appends `lastConfirmedAt` and confirmation provenance while retaining the original fact. + +Every ratification checks the approved payload digest and rechecks active facts under the same lock as its append. A conflicting proposal cannot silently become current because another fact was approved first. Repeated ratification of one proposal cannot append twice. The CLI durably claims review items and reports failed ratifications with exit 5. Legacy approvals without digests need a fresh proposal and approval. diff --git a/docs/releasing.md b/docs/releasing.md index 306f2e3..f1f490b 100644 --- a/docs/releasing.md +++ b/docs/releasing.md @@ -14,9 +14,13 @@ npm ci npm test npm run check:publish npm run check:release -mcp-publisher validate server.json +# Validate server.json against its declared JSON schema without publishing. ``` +Some publisher versions advertise `validate` but do not implement it. In that +case use a JSON Schema validator against the schema declared in `server.json`; +do not use `publish` as a validation probe. + Check npm credentials with `npm whoami`. Publish core first, then the public adapters, then CLI and MCP; consumers must not receive a package whose required ATS dependency version is missing from npm. Use `npm publish --access public diff --git a/package-lock.json b/package-lock.json index b64afe6..2f0c681 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "agentic-task-system", - "version": "0.15.0", + "version": "0.16.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "agentic-task-system", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "workspaces": [ "packages/*" @@ -2127,98 +2127,98 @@ }, "packages/adapter-airtable": { "name": "@reneza/ats-adapter-airtable", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-beads": { "name": "@reneza/ats-adapter-beads", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-composite": { "name": "@reneza/ats-adapter-composite", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-github": { "name": "@reneza/ats-adapter-github", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-google": { "name": "@reneza/ats-adapter-google", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-notion": { "name": "@reneza/ats-adapter-notion", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-obsidian": { "name": "@reneza/ats-adapter-obsidian", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-okf": { "name": "@reneza/ats-adapter-okf", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-taskmaster": { "name": "@reneza/ats-adapter-taskmaster", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-ticktick": { "name": "@reneza/ats-adapter-ticktick", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/adapter-ticktick-cache": { "name": "@reneza/ats-adapter-ticktick-cache", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } }, "packages/cli": { "name": "@reneza/ats-cli", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "dependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "bin": { "ats": "bin/ats.js" @@ -2226,23 +2226,23 @@ }, "packages/core": { "name": "@reneza/ats-core", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT" }, "packages/mcp": { "name": "@reneza/ats-mcp", - "version": "0.15.0", + "version": "0.16.0", "license": "MIT", "dependencies": { "@modelcontextprotocol/sdk": "^1.30.1", - "@reneza/ats-core": "^0.15.0", + "@reneza/ats-core": "^0.16.0", "zod": "^4.6.5" }, "bin": { "ats-mcp": "server.js" }, "peerDependencies": { - "@reneza/ats-adapter-ticktick": "^0.15.0" + "@reneza/ats-adapter-ticktick": "^0.16.0" }, "peerDependenciesMeta": { "@reneza/ats-adapter-ticktick": { diff --git a/package.json b/package.json index dd3ce61..59b9751 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "agentic-task-system", - "version": "0.15.0", + "version": "0.16.0", "private": true, "type": "module", "description": "Your task manager is the best agent memory you're not using. Agent-native context layer over your existing task app, with hybrid retrieval (RRF) and pluggable storage adapters.", diff --git a/packages/adapter-airtable/package.json b/packages/adapter-airtable/package.json index e5e5d33..23b697e 100644 --- a/packages/adapter-airtable/package.json +++ b/packages/adapter-airtable/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-airtable", - "version": "0.15.0", + "version": "0.16.0", "description": "Airtable adapter for Agentic Task System. Expose any Airtable base as agent-queryable records through ATS retrieval, RRF fusion, and MCP — a table is a project, a record is a task. Adapter, not migration.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "test": "node --test" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-beads/package.json b/packages/adapter-beads/package.json index 6415596..b4ab368 100644 --- a/packages/adapter-beads/package.json +++ b/packages/adapter-beads/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-beads", - "version": "0.15.0", + "version": "0.16.0", "description": "Beads adapter for Agentic Task System using the official bd JSON CLI over repository-local Dolt state.", "type": "module", "main": "index.js", @@ -13,7 +13,7 @@ "prepublishOnly": "node ../../scripts/check-no-pii.mjs --self" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "repository": { "type": "git", diff --git a/packages/adapter-composite/package.json b/packages/adapter-composite/package.json index 0cff042..722bdd3 100644 --- a/packages/adapter-composite/package.json +++ b/packages/adapter-composite/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-composite", - "version": "0.15.0", + "version": "0.16.0", "description": "Cross-source adapter for the Agentic Task System. Query GitHub + Notion + TickTick + any ATS backends as ONE fused corpus: a single `ats find` returns one RRF-ranked list across all of them, each result tagged with its backend. The thing a single-vendor MCP server can't do.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "test": "node --test" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-github/api.js b/packages/adapter-github/api.js index 55623d5..7f522a8 100644 --- a/packages/adapter-github/api.js +++ b/packages/adapter-github/api.js @@ -205,8 +205,31 @@ export async function patchIssue(owner, repo, number, body, cfg = loadConfig()) export async function searchIssues(query, cfg = loadConfig()) { const scope = cfg.repos.map((r) => `repo:${r}`).join(' '); const q = `${query} is:issue${scope ? ' ' + scope : ''}`.trim(); - const { json } = await gh('/search/issues', { query: { q, per_page: PER_PAGE }, cfg }); - return Array.isArray(json.items) ? json.items.filter((it) => !isPullRequest(it)) : []; + const items = []; + const warnings = []; + const seen = new Set(); + const allowed = new Set(cfg.repos.map((repo) => repo.toLowerCase())); + for (let page = 1; page <= 10; page++) { + const { json, res } = await gh('/search/issues', { query: { q, per_page: PER_PAGE, page }, cfg }); + if (!Array.isArray(json.items)) throw new Error('GitHub search: invalid items response'); + if (json.incomplete_results) warnings.push({ source: 'github', error: 'GitHub returned incomplete search results' }); + for (const item of json.items) { + const repo = String(item.repository_url || '').match(/repos\/([^/]+\/[^/]+)\/?$/)?.[1]; + if (isPullRequest(item) || (allowed.size && (!repo || !allowed.has(repo.toLowerCase())))) continue; + const key = `${repo}:${item.number}`; + if (!seen.has(key)) { seen.add(key); items.push(item); } + } + if (page === 10 && (json.total_count > 1000 || nextLink(res))) { + warnings.push({ source: 'github', error: 'GitHub search is capped at 1000 matches; narrow the query' }); + } + if (!nextLink(res) && !(json.total_count > page * PER_PAGE)) break; + if (json.items.length === 0) { + warnings.push({ source: 'github', error: 'GitHub search stopped before the declared match count' }); + break; + } + } + items.warnings = warnings; + return items; } export function urlForIssue(owner, repo, number) { diff --git a/packages/adapter-github/index.js b/packages/adapter-github/index.js index 8d6b2e1..f5b4071 100644 --- a/packages/adapter-github/index.js +++ b/packages/adapter-github/index.js @@ -91,18 +91,25 @@ const adapter = { const cfg = loadConfig(); const repos = await listRepos(cfg); const tasks = []; + adapter.__fetchWarnings = []; for (const r of repos) { const { owner, repo } = splitProjectId(r.id); - const issues = await listIssues(owner, repo, cfg); - for (const it of issues) tasks.push(issueToTask(owner, repo, it)); + try { + const issues = await listIssues(owner, repo, cfg); + for (const it of issues) tasks.push(issueToTask(owner, repo, it)); + } catch (error) { + adapter.__fetchWarnings.push({ source: r.id, error: error.message }); + } } return tasks; }, async searchByQuery(query) { const q = String(query || '').trim(); + adapter.__searchWarnings = []; if (!q) return []; const items = await searchIssues(q, loadConfig()); + adapter.__searchWarnings = items.warnings || []; return items.map((it) => { const repoUrl = String(it.repository_url || ''); const m = repoUrl.match(/repos\/([^/]+)\/([^/]+)\/?$/); diff --git a/packages/adapter-github/package.json b/packages/adapter-github/package.json index 9a6d97d..7b5dc4b 100644 --- a/packages/adapter-github/package.json +++ b/packages/adapter-github/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-github", - "version": "0.15.0", + "version": "0.16.0", "description": "GitHub Issues adapter for Agentic Task System. Expose any repository's issues as agent-queryable tasks through ATS retrieval, RRF fusion, and MCP — a repo is a project, an issue is a task. Adapter, not migration.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "test": "node --test" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-github/test/github.test.js b/packages/adapter-github/test/github.test.js index 2144b1c..49e20e4 100644 --- a/packages/adapter-github/test/github.test.js +++ b/packages/adapter-github/test/github.test.js @@ -199,3 +199,44 @@ test('authStatus reports authenticated with the login from /user', async () => { assert.equal(s.authenticated, true); assert.equal(s.login, 'octo'); }); + + +test('native search follows pages, deduplicates, and preserves configured repository scope', async () => { + process.env.ATS_GITHUB_REPOS = 'octo/widgets'; + try { + globalThis.fetch = async (url) => { + const page = Number(new URL(url).searchParams.get('page')); + const item = page === 1 ? ISSUE1 : ISSUE2; + return new globalThis.Response(JSON.stringify({ total_count: 101, items: [ + { ...item, repository_url: 'https://api.github.com/repos/octo/widgets' }, + { ...ISSUE1, number: 999, repository_url: 'https://api.github.com/repos/other/private' }, + ] }), { headers: page === 1 ? { link: '; rel="next"' } : {} }); + }; + const tasks = await adapter.searchByQuery('invoice repo:other/private'); + assert.deepEqual(tasks.map((task) => task.id), ['7', '8']); + assert.deepEqual(adapter.__searchWarnings, []); + } finally { delete process.env.ATS_GITHUB_REPOS; } +}); + +test('native search reports GitHub server incompleteness and the 1000 result ceiling', async () => { + globalThis.fetch = async (url) => { + const page = Number(new URL(url).searchParams.get('page')); + return new globalThis.Response(JSON.stringify({ total_count: 1200, incomplete_results: true, items: [{ ...ISSUE1, number: page, repository_url: 'https://api.github.com/repos/octo/widgets' }] })); + }; + assert.equal((await adapter.searchByQuery('common')).length, 10); + assert.ok(adapter.__searchWarnings.some((warning) => /1000/.test(warning.error))); + assert.ok(adapter.__searchWarnings.some((warning) => /incomplete/.test(warning.error))); +}); + +test('bulk corpus preserves a healthy repository and reports a failed repository', async () => { + process.env.ATS_GITHUB_REPOS = 'octo/widgets,other/unavailable'; + const healthyFetch = globalThis.fetch; + try { + globalThis.fetch = async (url, init) => new URL(url).pathname.startsWith('/repos/other/') + ? new globalThis.Response('{"message":"repository unavailable"}', { status: 404 }) : healthyFetch(url, init); + const tasks = await adapter.bulkFetch(); + assert.equal(tasks.length, 2); + assert.equal(adapter.__fetchWarnings.length, 1); + assert.equal(adapter.__fetchWarnings[0].source, 'other/unavailable'); + } finally { delete process.env.ATS_GITHUB_REPOS; } +}); diff --git a/packages/adapter-google/package.json b/packages/adapter-google/package.json index 4b78bcc..adbb12a 100644 --- a/packages/adapter-google/package.json +++ b/packages/adapter-google/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-google", - "version": "0.15.0", + "version": "0.16.0", "description": "Google Workspace adapter for Agentic Task System. Pull Google Sheets, Docs, and Slides into ATS retrieval and MCP as a read-only corpus, authed as a dedicated share-scoped user. Adapter, not migration.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "test": "node --test" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-notion/api.js b/packages/adapter-notion/api.js index 6bebe63..f459aa3 100644 --- a/packages/adapter-notion/api.js +++ b/packages/adapter-notion/api.js @@ -62,6 +62,7 @@ export async function notion(apiPath, opts = {}) { const url = new URL(cfg.endpoint.replace(/\/$/, '') + apiPath); const init = { method: opts.method || 'GET', + retrySafe: apiPath === '/v1/search' || /^\/v1\/databases\/[^/]+\/query$/.test(apiPath), headers: { Authorization: `Bearer ${cfg.token}`, 'Notion-Version': NOTION_VERSION, diff --git a/packages/adapter-notion/index.js b/packages/adapter-notion/index.js index 621fe57..ed314db 100644 --- a/packages/adapter-notion/index.js +++ b/packages/adapter-notion/index.js @@ -92,25 +92,50 @@ const adapter = { const cfg = loadConfig(); const dbs = await listDatabases(cfg); const tasks = []; + adapter.__fetchWarnings = []; for (const db of dbs) { - const pages = await queryDatabase(db.id, cfg); - for (const p of pages) tasks.push(pageToTask(p)); + try { + const pages = await queryDatabase(db.id, cfg); + for (const p of pages) tasks.push(pageToTask(p)); + } catch (error) { + adapter.__fetchWarnings.push({ source: db.id, error: error.message }); + } } return tasks; }, async searchByQuery(query) { const q = String(query || '').trim(); + adapter.__searchWarnings = []; if (!q) return []; const cfg = loadConfig(); - const res = await notion('/v1/search', { - method: 'POST', - body: { query: q, filter: { property: 'object', value: 'page' }, page_size: 100 }, - cfg, - }); - return (res.results || []) - .filter((r) => r.object === 'page' && r.parent?.database_id) - .map((p) => pageToTask(p)); + const normalizeId = (id) => String(id).replace(/-/g, '').toLowerCase(); + const allowed = new Set(cfg.databases.map(normalizeId)); + const tasks = new Map(); + const cursors = new Set(); + let cursor; + for (let page = 0; page < 100; page++) { + const res = await notion('/v1/search', { + method: 'POST', + body: { query: q, filter: { property: 'object', value: 'page' }, page_size: 100, start_cursor: cursor }, + cfg, + }); + if (!Array.isArray(res.results)) throw new Error('Notion search: invalid results response'); + for (const record of res.results) { + if (record.object !== 'page' || !record.parent?.database_id) continue; + if (allowed.size && !allowed.has(normalizeId(record.parent.database_id))) continue; + tasks.set(record.id, pageToTask(record)); + } + if (!res.has_more) return [...tasks.values()]; + cursor = res.next_cursor; + if (!cursor || cursors.has(cursor)) { + adapter.__searchWarnings.push({ source: 'notion', error: 'Notion search returned a missing or repeated cursor' }); + return [...tasks.values()]; + } + cursors.add(cursor); + } + adapter.__searchWarnings.push({ source: 'notion', error: 'Notion search reached the 100-page bound; narrow the query' }); + return [...tasks.values()]; }, // ---- auth lifecycle -------------------------------------------------------- diff --git a/packages/adapter-notion/package.json b/packages/adapter-notion/package.json index b07941d..4fd99ac 100644 --- a/packages/adapter-notion/package.json +++ b/packages/adapter-notion/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-notion", - "version": "0.15.0", + "version": "0.16.0", "description": "Notion adapter for Agentic Task System. Expose any Notion database as agent-queryable pages through ATS retrieval, RRF fusion, and MCP — a database is a project, a page is a task. Adapter, not migration.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "test": "node --test" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-notion/test/notion.test.js b/packages/adapter-notion/test/notion.test.js index a01c5e0..f8d5a8e 100644 --- a/packages/adapter-notion/test/notion.test.js +++ b/packages/adapter-notion/test/notion.test.js @@ -194,3 +194,33 @@ test('authStatus reports authenticated and counts databases', async () => { assert.equal(s.authenticated, true); assert.equal(s.databases, 1); }); + + +test('native search paginates, applies configured database scope and warns about repeated cursors', async () => { + process.env.ATS_NOTION_DATABASES = 'db-1111'; + let pages = 0; + try { + globalThis.fetch = async (_url, init) => { + const input = JSON.parse(init.body); + const next = input.start_cursor ? { ...PAGE2, id: 'third' } : PAGE1; + pages++; + return new globalThis.Response(JSON.stringify({ results: [next, { ...PAGE2, id: 'outside', parent: { database_id: 'other-db' } }], has_more: pages < 3, next_cursor: 'cursor-1' })); + }; + const tasks = await adapter.searchByQuery('invoice'); + assert.deepEqual(tasks.map((task) => task.id), [PAGE1.id, 'third']); + assert.equal(pages, 2); + assert.match(adapter.__searchWarnings[0].error, /repeated cursor/); + } finally { delete process.env.ATS_NOTION_DATABASES; } +}); + +test('bulk corpus preserves a healthy database and reports a failed database', async () => { + process.env.ATS_NOTION_DATABASES = 'db-1111,missing'; + const healthyFetch = globalThis.fetch; + try { + globalThis.fetch = async (url, init) => new URL(url).pathname.includes('/missing') + ? new globalThis.Response('{"message":"database unavailable"}', { status: 404 }) : healthyFetch(url, init); + assert.equal((await adapter.bulkFetch()).length, 2); + assert.equal(adapter.__fetchWarnings.length, 1); + assert.equal(adapter.__fetchWarnings[0].source, 'missing'); + } finally { delete process.env.ATS_NOTION_DATABASES; } +}); diff --git a/packages/adapter-obsidian/package.json b/packages/adapter-obsidian/package.json index 5e1f9d8..b165503 100644 --- a/packages/adapter-obsidian/package.json +++ b/packages/adapter-obsidian/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-obsidian", - "version": "0.15.0", + "version": "0.16.0", "description": "Obsidian vault adapter for Agentic Task System. Hybrid + RRF retrieval, the wiki layer, and the MCP server over the folder of markdown you already keep in Obsidian — adapter, not migration.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "prepublishOnly": "node ../../scripts/check-no-pii.mjs --self" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-okf/package.json b/packages/adapter-okf/package.json index 38c05cd..81fea2f 100644 --- a/packages/adapter-okf/package.json +++ b/packages/adapter-okf/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-okf", - "version": "0.15.0", + "version": "0.16.0", "description": "OKF bundle adapter for Agentic Task System. Expose Open Knowledge Format markdown bundles through ATS retrieval, graph context, and MCP.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "prepublishOnly": "node ../../scripts/check-no-pii.mjs --self" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-taskmaster/package.json b/packages/adapter-taskmaster/package.json index fad2476..414c0d5 100644 --- a/packages/adapter-taskmaster/package.json +++ b/packages/adapter-taskmaster/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-taskmaster", - "version": "0.15.0", + "version": "0.16.0", "description": "Local Taskmaster adapter for Agentic Task System: cross-tag search, task context, and field-preserving writes over tasks.json.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "prepublishOnly": "node ../../scripts/check-no-pii.mjs --self" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/adapter-ticktick-cache/package.json b/packages/adapter-ticktick-cache/package.json index 881f7e8..8915038 100644 --- a/packages/adapter-ticktick-cache/package.json +++ b/packages/adapter-ticktick-cache/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-ticktick-cache", - "version": "0.15.0", + "version": "0.16.0", "private": true, "description": "Local-first ATS adapter over a centralized TickTick JSON cache, with OpenAPI-backed synchronization and writes.", "type": "module", @@ -10,6 +10,6 @@ ], "license": "MIT", "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" } } diff --git a/packages/adapter-ticktick/package.json b/packages/adapter-ticktick/package.json index b6a07c4..f50112b 100644 --- a/packages/adapter-ticktick/package.json +++ b/packages/adapter-ticktick/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-adapter-ticktick", - "version": "0.15.0", + "version": "0.16.0", "description": "Reference TickTick adapter for Agentic Task System. Wraps TickTick OpenAPI v1 + qdrant + ollama (nomic-embed) into the ATS adapter contract.", "type": "module", "main": "index.js", @@ -12,7 +12,7 @@ "prepublishOnly": "node ../../scripts/check-no-pii.mjs --self" }, "peerDependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "dependencies": {}, "repository": { diff --git a/packages/cli/bin/ats.js b/packages/cli/bin/ats.js index 2bc3e6c..6bad05a 100755 --- a/packages/cli/bin/ats.js +++ b/packages/cli/bin/ats.js @@ -38,7 +38,7 @@ import { getAgentLayerHelp, } from '../parser.js'; import { formatSkipped } from '../format-skip.js'; -import { lookupIdempotencyKey, recordIdempotencyKey, findActiveByTitle } from '../idempotency.js'; +import { recordIdempotencyKey, claimIdempotencyKey, claimReviewedIdempotencyKey, markIdempotencyUncertain, findActiveByTitle } from '../idempotency.js'; import { validateAdapter, runConformance, @@ -94,6 +94,7 @@ import { markReviewItemApplied, claimReviewItem, reviewTargetRevision, + stableDigest, exportState, importState, inspectState, @@ -103,6 +104,8 @@ import { loadFacts, proposeFact, proposeRetract, + proposeConfirm, + staleFacts, ratifyFactItem, listKgFacts, askFacts, @@ -124,8 +127,8 @@ import { resolveOpen, formatOpenResult, launchUrl, shouldLaunch } from '../open. import { readStructuredInput, parseBatchInput, - readBatchJournal, - appendBatchJournal, + claimBatchItem, + finishBatchItem, withTimeout, classifyError, errorEnvelope, @@ -1107,23 +1110,25 @@ async function applyReviewedWrite(item, adapter, t) { return result; } case 'task.created': { - // A staged idempotent create that already landed (an earlier apply of the - // same key) is not created twice. + let claim; if (p.idempotencyKey !== undefined) { - const seen = lookupIdempotencyKey(p.idempotencyKey, { configDir: atsConfigDir() }); - if (seen?.taskId) return { created: false, idempotent: true, key: p.idempotencyKey, task: seen }; + if (p.idempotencyScope !== process.env.ATS_CORPUS_SCOPE) throw withExitCode(new Error('Idempotency precondition failed: active source changed since staging.'), 3); + claim = claimReviewedIdempotencyKey(p.idempotencyKey, item.id, p.idempotencyDigest, { configDir: atsConfigDir() }); + if (claim.replayed) return { created: false, idempotent: true, key: p.idempotencyKey, task: claim.entry }; + } + const idemOptions = { configDir: atsConfigDir(), token: claim?.entry.token }; + try { + const result = t?.create + ? await t.create(p.projectId || '', p.title, p.opts || {}) + : await adapter.createTask({ title: p.title, projectId: p.projectId || undefined, + content: p.opts?.content, dueDate: p.opts?.dueDate, tags: tagsToArray(p.opts?.tags) }); + auditCliWrite('task.created', result, { projectId: p.projectId }, { title: p.title, reviewId: item.id }, false, undefined, approvals); + if (claim) recordIdempotencyKey(p.idempotencyKey, taskRef(result, p.projectId), idemOptions); + return result; + } catch (error) { + if (claim) { try { markIdempotencyUncertain(p.idempotencyKey, idemOptions); } catch {} } + throw error; } - const result = t?.create - ? await t.create(p.projectId || '', p.title, p.opts || {}) - : await adapter.createTask({ - title: p.title, - projectId: p.projectId || undefined, - content: p.opts?.content, - dueDate: p.opts?.dueDate, - tags: tagsToArray(p.opts?.tags), - }); - auditCliWrite('task.created', result, { projectId: p.projectId }, { title: p.title, reviewId: item.id }, false, undefined, approvals); - return result; } default: throw new Error(`Unknown staged write action: ${p.action}`); @@ -1154,7 +1159,7 @@ async function handleKg() { const text = file === '-' ? fs.readFileSync(0, 'utf8') : fs.readFileSync(file, 'utf8'); const report = proposeFactLines(text.split('\n'), { by: agentId, - defaults: { domain: args.options.domain, source: args.options.source, confidence: args.options.confidence }, + defaults: { domain: args.options.domain, source: args.options.source, confidence: args.options.confidence, validAt: args.options['valid-at'], learnedAt: args.options['learned-at'] }, }); const ok = report.refused === 0 && report.invalid === 0; report.message = ok @@ -1179,6 +1184,8 @@ async function handleKg() { try { item = proposeFact({ subject, predicate, object, + validAt: args.options['valid-at'], + learnedAt: args.options['learned-at'], domain: args.options.domain, source: args.options.source, confidence: args.options.confidence, @@ -1199,6 +1206,12 @@ async function handleKg() { message: `Fact proposed as ${item.id.slice(0, 8)}. Ratify with: ats review approve ${item.id.slice(0, 8)} && ats kg ratify --all`, }; } + case 'stale': + return staleFacts({ domain: args.options.domain, days: args.options.days === undefined ? 60 : Number(args.options.days), limit: args.options.limit === undefined ? 50 : Number(args.options.limit) }); + case 'confirm': { + const item = proposeConfirm({ factId: args.positional[0], source: args.options.source, by: agentId }); + return { staged: true, reviewId: item.id, message: 'Confirmation staged; approve it and run ats kg ratify.' }; + } case 'retract': { if (!args.positional[0]) { console.error('Usage: ats kg retract FACT_ID [--reason "..."]'); process.exit(1); } const item = proposeRetract({ factId: args.positional[0], reason: args.options.reason, by: agentId }); @@ -1219,17 +1232,20 @@ async function handleKg() { process.exit(1); } const ratified = []; - for (const item of targets) { + for (const target of targets) { + let item; try { + item = claimReviewItem(target.id); const outcome = ratifyFactItem(item); const result = outcome.op === 'add' ? { factId: outcome.fact.id, ...(outcome.superseded ? { superseded: outcome.superseded } : {}) } - : { retracted: outcome.factId }; - markReviewItemApplied(item.id, { result }); + : { [outcome.op === 'confirm' ? 'confirmed' : 'retracted']: outcome.factId }; + markReviewItemApplied(item.id, { result, applyToken: item.applyToken }); ratified.push({ id: item.id.slice(0, 8), ok: true, ...result }); } catch (err) { - try { markReviewItemApplied(item.id, { error: err.message }); } catch {} - ratified.push({ id: item.id.slice(0, 8), ok: false, error: err.message }); + if (item) { try { markReviewItemApplied(item.id, { error: err.message, applyToken: item.applyToken }); } catch {} } + ratified.push({ id: target.id.slice(0, 8), ok: false, error: err.message }); + process.exitCode = 5; } } return { ratified }; @@ -1441,25 +1457,30 @@ async function handleBatch() { const file = args.subcommand || args.positional[0]; const items = parseBatchInput(file); const journal = typeof args.options.journal === 'string' ? args.options.journal : null; - const completed = readBatchJournal(journal); const dryRun = args.options['dry-run'] === true; const adapter = await loadAdapter(); const t = adapter.__ext?.tasks; const outcomes = []; for (const item of items) { - if (completed.has(item.id)) { - outcomes.push({ id: item.id, op: item.op, status: 'skipped', reason: 'already applied in journal' }); - continue; - } let outcome; + let claim; try { + if (!dryRun) claim = claimBatchItem(journal, item, process.env.ATS_CORPUS_SCOPE); + if (claim?.replayed) { + outcomes.push({ id: item.id, op: item.op, status: 'skipped', reason: 'matching operation already applied in journal' }); + continue; + } const executed = await executeBatchOperation(adapter, t, item, dryRun); outcome = { id: item.id, op: item.op, ...executed }; } catch (error) { outcome = { id: item.id, op: item.op, status: 'failed', error: errorEnvelope(error).error }; } + // A dry run neither claims nor updates a real resume journal. + if (claim) { + try { finishBatchItem(journal, claim, outcome); } + catch (error) { outcome = { id: item.id, op: item.op, status: 'failed', error: errorEnvelope(error).error }; } + } outcomes.push(outcome); - appendBatchJournal(journal, outcome); } const summary = Object.fromEntries(['planned', 'applied', 'staged', 'skipped', 'failed'].map((status) => [status, outcomes.filter((outcome) => outcome.status === status).length])); if (summary.failed) process.exitCode = 5; @@ -1779,66 +1800,73 @@ async function handleTasks() { // Idempotent creates: a key that already produced something returns it; // --if-absent returns the active task that already carries this title. const idemKey = typeof args.options['idempotency-key'] === 'string' ? args.options['idempotency-key'] : undefined; - if (idemKey !== undefined) { - const seen = lookupIdempotencyKey(idemKey, { configDir: atsConfigDir() }); - if (seen) return await idempotentReplay(idemKey, seen, adapter, t); - } - if (args.options['if-absent'] === true) { - const existing = await activeTaskWithTitle(adapter, t, projectId, title); - if (existing) { - if (idemKey !== undefined) recordIdempotencyKey(idemKey, taskRef(existing, projectId), { configDir: atsConfigDir() }); - return { - created: false, - existing: true, - reason: 'if-absent: an active task with this title already exists', - task: existing, - }; + const binding = stableDigest({ scope: process.env.ATS_CORPUS_SCOPE, request: { projectId, title, opts, ifAbsent: args.options['if-absent'] === true } }); + const claim = idemKey === undefined ? null : claimIdempotencyKey(idemKey, + { projectId, title, opts, ifAbsent: args.options['if-absent'] === true }, + { configDir: atsConfigDir(), scope: process.env.ATS_CORPUS_SCOPE }); + if (claim?.replayed) return await idempotentReplay(idemKey, claim.entry, adapter, t); + const idemOptions = { configDir: atsConfigDir(), ...(claim ? { token: claim.entry.token } : {}) }; + try { + if (args.options['if-absent'] === true) { + const existing = await activeTaskWithTitle(adapter, t, projectId, title); + if (existing) { + if (idemKey !== undefined) recordIdempotencyKey(idemKey, taskRef(existing, projectId), idemOptions); + return { + created: false, + existing: true, + reason: 'if-absent: an active task with this title already exists', + task: existing, + }; + } } - } - // Creates have no target metadata to consult; they stage only under the - // global ATS_REVIEW_ALL=1 gate. - const createGate = reviewGate('task.created', null, { projectId, title, opts, ...(idemKey !== undefined ? { idempotencyKey: idemKey } : {}) }); - if (createGate) { - if (idemKey !== undefined) recordIdempotencyKey(idemKey, { reviewId: createGate.reviewId, projectId }, { configDir: atsConfigDir() }); - return createGate; - } - const result = t?.create - ? await t.create(projectId, title, opts) - : await adapter.createTask({ - title, - projectId: projectId || undefined, - content: opts.content, - dueDate: opts.dueDate, - tags: tagsToArray(opts.tags), - }); - const action = auditCliWrite('task.created', result, { projectId }, { title, ...(idemKey !== undefined ? { idempotencyKey: idemKey } : {}) }); - if (idemKey !== undefined) recordIdempotencyKey(idemKey, taskRef(result, projectId), { configDir: atsConfigDir() }); - const relevance = adapter.__ext?.relevance; - if (relevance?.isEnabled?.({ - relevance: !!args.options.relevance, - noRelevance: args.options['no-relevance'] === true, - })) { - try { - const block = await relevance.buildEnrichInstruction({ - taskId: result.task?.fullId || result.task?.id, - projectId: result.task?.fullProjectId || result.task?.projectId || projectId, - title: result.task?.title || title, - content: opts.content || '', - wikiProject: wikiProject(), + // Creates have no target metadata to consult; they stage only under the + // global ATS_REVIEW_ALL=1 gate. + const createGate = reviewGate('task.created', null, { projectId, title, opts, ...(idemKey !== undefined ? { idempotencyKey: idemKey, idempotencyDigest: binding, idempotencyScope: process.env.ATS_CORPUS_SCOPE } : {}) }); + if (createGate) { + if (idemKey !== undefined) recordIdempotencyKey(idemKey, { reviewId: createGate.reviewId, projectId }, idemOptions); + return createGate; + } + const result = t?.create + ? await t.create(projectId, title, opts) + : await adapter.createTask({ + title, + projectId: projectId || undefined, + content: opts.content, + dueDate: opts.dueDate, + tags: tagsToArray(opts.tags), }); - if (block) result._relevanceInstruction = block; - } catch (err) { - console.error(`Warning: relevance enrichment failed: ${err.message}`); + const action = auditCliWrite('task.created', result, { projectId }, { title, ...(idemKey !== undefined ? { idempotencyKey: idemKey, idempotencyDigest: binding, idempotencyScope: process.env.ATS_CORPUS_SCOPE } : {}) }); + if (idemKey !== undefined) recordIdempotencyKey(idemKey, taskRef(result, projectId), idemOptions); + const relevance = adapter.__ext?.relevance; + if (relevance?.isEnabled?.({ + relevance: !!args.options.relevance, + noRelevance: args.options['no-relevance'] === true, + })) { + try { + const block = await relevance.buildEnrichInstruction({ + taskId: result.task?.fullId || result.task?.id, + projectId: result.task?.fullProjectId || result.task?.projectId || projectId, + title: result.task?.title || title, + content: opts.content || '', + wikiProject: wikiProject(), + }); + if (block) result._relevanceInstruction = block; + } catch (err) { + console.error(`Warning: relevance enrichment failed: ${err.message}`); + } } + return attachMutationReceipt(result, mutationReceipt({ + operation: 'create', + action, + requested: { projectId: projectId || null, title, ...opts }, + observed: snapshotTask(result), + verified: taskValue(result).title === title, + target: taskRefFromResult(result, { projectId }), + })); + } catch (error) { + if (claim) { try { markIdempotencyUncertain(idemKey, idemOptions); } catch {} } + throw error; } - return attachMutationReceipt(result, mutationReceipt({ - operation: 'create', - action, - requested: { projectId: projectId || null, title, ...opts }, - observed: snapshotTask(result), - verified: taskValue(result).title === title, - target: taskRefFromResult(result, { projectId }), - })); } case 'update': { const input = args.options.input diff --git a/packages/cli/idempotency.js b/packages/cli/idempotency.js index 91fecb9..c252aa3 100644 --- a/packages/cli/idempotency.js +++ b/packages/cli/idempotency.js @@ -8,7 +8,8 @@ // (default 7 days). import fs from 'node:fs'; import path from 'node:path'; -import { withLockSync, writeFileAtomicSync } from '@reneza/ats-core'; +import { randomUUID } from 'node:crypto'; +import { withLockSync, writeFileAtomicSync, stableDigest } from '@reneza/ats-core'; export const IDEMPOTENCY_VERSION = 1; const TTL_MS = Number(process.env.ATS_IDEMPOTENCY_TTL_MS) || 7 * 24 * 60 * 60 * 1000; @@ -18,11 +19,61 @@ export function idempotencyPath(configDir) { } function readStore(file) { - try { - const raw = JSON.parse(fs.readFileSync(file, 'utf8')); - if (raw && raw.version === IDEMPOTENCY_VERSION && raw.keys && typeof raw.keys === 'object') return raw; - } catch {} - return { version: IDEMPOTENCY_VERSION, keys: {} }; + if (!fs.existsSync(file)) return { version: IDEMPOTENCY_VERSION, keys: Object.create(null) }; + const raw = JSON.parse(fs.readFileSync(file, 'utf8')); + if (!raw || raw.version !== IDEMPOTENCY_VERSION || !raw.keys || typeof raw.keys !== 'object' || Array.isArray(raw.keys)) throw new Error('Invalid idempotency store; inspect it before retrying writes.'); + raw.keys = Object.assign(Object.create(null), raw.keys); + return raw; +} + +function conflict(message) { + return Object.assign(new Error(`Idempotency precondition failed: ${message}`), { code: 'ATS_PRECONDITION', exitCode: 3 }); +} + +/** Bind a durable pre-write claim to the exact request and active source. */ +export function claimIdempotencyKey(key, request, { configDir, scope, now = Date.now() } = {}) { + const file = idempotencyPath(configDir); + const digest = stableDigest({ scope, request }); + return withLockSync(file, () => { + const store = readStore(file); + const prior = store.keys[key]; + // An uncertain outcome never expires into an automatic duplicate write. + if (prior && (live(prior, now) || ['claimed', 'uncertain'].includes(prior.status))) { + if (prior.digest !== digest) throw conflict('key has no matching request/source binding; use a new key for different work.'); + if (['claimed', 'uncertain'].includes(prior.status)) throw conflict('previous outcome is in flight or uncertain; inspect the backend before staging fresh work.'); + return { replayed: true, entry: prior }; + } + const entry = { digest, status: 'claimed', token: randomUUID(), at: now }; + store.keys[key] = entry; + writeFileAtomicSync(file, JSON.stringify(store, null, 2) + '\n'); + return { replayed: false, entry }; + }, { label: 'idempotency keys' }); +} + +/** Move this staged create's key into an applying claim under the same store lock. */ +export function claimReviewedIdempotencyKey(key, reviewId, digest, { configDir } = {}) { + const file = idempotencyPath(configDir); + return withLockSync(file, () => { + const store = readStore(file); + const prior = store.keys[key]; + if (!prior || prior.digest !== digest) throw conflict('reviewed create has no matching stored binding.'); + if (['claimed', 'uncertain'].includes(prior.status)) throw conflict('previous outcome is in flight or uncertain.'); + if (prior.taskId) return { replayed: true, entry: prior }; + if (prior.reviewId !== reviewId) throw conflict('key belongs to a different review item.'); + store.keys[key] = { ...prior, status: 'claimed', token: randomUUID(), at: Date.now() }; + writeFileAtomicSync(file, JSON.stringify(store, null, 2) + '\n'); + return { replayed: false, entry: store.keys[key] }; + }, { label: 'idempotency keys' }); +} + +export function markIdempotencyUncertain(key, { configDir, token } = {}) { + const file = idempotencyPath(configDir); + return withLockSync(file, () => { + const store = readStore(file); + if (store.keys[key]?.token !== token) throw conflict('claim token does not match.'); + store.keys[key].status = 'uncertain'; + writeFileAtomicSync(file, JSON.stringify(store, null, 2) + '\n'); + }, { label: 'idempotency keys' }); } function live(entry, now) { @@ -36,15 +87,17 @@ export function lookupIdempotencyKey(key, { configDir, now = Date.now() } = {}) } /** Record what a key produced. Expired keys are pruned on every write. */ -export function recordIdempotencyKey(key, outcome, { configDir, now = Date.now() } = {}) { +export function recordIdempotencyKey(key, outcome, { configDir, now = Date.now(), token } = {}) { const file = idempotencyPath(configDir); fs.mkdirSync(path.dirname(file), { recursive: true, mode: 0o700 }); return withLockSync(file, () => { const store = readStore(file); + if (token && store.keys[key]?.token !== token) throw conflict('claim token does not match.'); + const prior = store.keys[key]; for (const [k, v] of Object.entries(store.keys)) { - if (!live(v, now)) delete store.keys[k]; + if (!live(v, now) && !['claimed', 'uncertain'].includes(v.status)) delete store.keys[k]; } - store.keys[key] = { ...outcome, at: now }; + store.keys[key] = { ...prior, ...outcome, status: 'done', at: now }; writeFileAtomicSync(file, JSON.stringify(store, null, 2) + '\n'); return store.keys[key]; }, { label: 'idempotency keys' }); diff --git a/packages/cli/package.json b/packages/cli/package.json index b18dfc6..6424b4d 100644 --- a/packages/cli/package.json +++ b/packages/cli/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-cli", - "version": "0.15.0", + "version": "0.16.0", "description": "Command-line interface for Agentic Task System. Routes to the configured storage adapter and exposes find / get / url / links / hybrid / similar / create / update / bench.", "type": "module", "bin": { @@ -15,7 +15,7 @@ "prepublishOnly": "node ../../scripts/check-no-pii.mjs --self" }, "dependencies": { - "@reneza/ats-core": "^0.15.0" + "@reneza/ats-core": "^0.16.0" }, "repository": { "type": "git", diff --git a/packages/cli/parser.js b/packages/cli/parser.js index 75fa339..ae32808 100644 --- a/packages/cli/parser.js +++ b/packages/cli/parser.js @@ -35,7 +35,7 @@ const VALUE_OPTIONS = new Set([ 'days', 'since', 'state', 'spool', 'interval', 'due-within-hours', 'max', 'threshold', 'max-corpus', 'rerank-depth', 'min-sources', 'facts-limit', 'depth', 'max-depth', 'max-nodes', 'restore', 'dir', 'keep', 'dupes', - 'status', 'stale-days', 'as-of', 'acknowledge-rejected', 'supersedes', 'object', 'dialect', 'center', 'outcome', 'done-when', 'parent-project', + 'valid-at', 'learned-at', 'status', 'stale-days', 'as-of', 'acknowledge-rejected', 'supersedes', 'object', 'dialect', 'center', 'outcome', 'done-when', 'parent-project', 'parent-task', 'approval-required', 'valid-from', 'valid-until', 'allow-actions', 'allow-resources', 'deny-resources', 'approval-actions', ]); @@ -788,9 +788,12 @@ every fact records who did both and from what source. Usage: ats kg propose SUBJ PRED OBJ [--domain D --source REF --confidence C --task P/T] [--supersedes FACT_ID | --additive] [--acknowledge-rejected ID] + [--valid-at ISO --learned-at ISO] ats kg propose --file FILE|- One JSON object per line, every line through the gate; --domain/--source/ --confidence fill what a line lacks + ats kg stale [--days N --domain D] Active facts due for evidence review + ats kg confirm FACT_ID --source REF Reviewed evidence confirmation ats kg retract FACT_ID [--reason "..."] Retraction proposal — reviewed too ats kg pending [--domain D] What each queued proposal would do, checked against the store now @@ -817,7 +820,9 @@ Usage: The gate: a duplicate of an active or queued fact is a no-op (exit 0); a contradiction (same subject + predicate, another object) needs --supersedes or --additive; a triple the reviewer declined needs --acknowledge-rejected. A -refusal prints its report and exits 4. +refusal prints its report and exits 4. Source timestamps use ISO timestamps +with a timezone and cannot name future or planned dates. Confirmation is +reviewed and retains the original validity interval and provenance. Fact proposals share the review queue: ats review list / approve / reject work on them (kind kg.fact). A retracted or superseded fact keeps its validity diff --git a/packages/cli/reliability.js b/packages/cli/reliability.js index 832eb3a..8377b3e 100644 --- a/packages/cli/reliability.js +++ b/packages/cli/reliability.js @@ -1,5 +1,7 @@ import fs from 'node:fs'; import path from 'node:path'; +import { randomUUID } from 'node:crypto'; +import { stableDigest, withLockSync } from '@reneza/ats-core'; export function readStructuredInput(file, { allowed = [], required = [], stdin = 0 } = {}) { if (!file || file === true) throw new Error('--input requires a JSON file path or - for stdin.'); @@ -47,12 +49,13 @@ export function parseBatchInput(file, { stdin = 0 } = {}) { } export function readBatchJournal(file) { - if (!file || !fs.existsSync(file)) return new Set(); - const completed = new Set(); + if (!file || !fs.existsSync(file)) return new Map(); + const completed = new Map(); for (const line of fs.readFileSync(file, 'utf8').split('\n').filter((value) => value.trim())) { try { const entry = JSON.parse(line); - if (['applied', 'staged'].includes(entry?.status) && entry.id) completed.add(entry.id); + if (!entry || typeof entry.id !== 'string' || typeof entry.status !== 'string') throw new Error('invalid journal entry'); + completed.set(entry.id, entry); } catch (error) { throw new Error(`Batch journal ${file} is invalid: ${error.message}`, { cause: error }); } @@ -60,10 +63,41 @@ export function readBatchJournal(file) { return completed; } -export function appendBatchJournal(file, outcome) { - if (!file) return; +function writeBatchJournal(file, outcome) { fs.mkdirSync(path.dirname(path.resolve(file)), { recursive: true, mode: 0o700 }); fs.appendFileSync(file, JSON.stringify({ ...outcome, recordedAt: new Date().toISOString() }) + '\n', { mode: 0o600 }); + fs.chmodSync(file, 0o600); +} + +export function appendBatchJournal(file, outcome) { + if (!file) return; + withLockSync(file, () => writeBatchJournal(file, outcome), { label: 'batch journal' }); +} + +/** Claim before a side effect; an unfinished claim is never automatically retried. */ +export function claimBatchItem(file, item, scope) { + if (!file) return null; + const digest = stableDigest({ scope, item }); + return withLockSync(file, () => { + const prior = readBatchJournal(file).get(item.id); + if (prior) { + if (prior.digest !== digest) throw new Error(`Batch precondition failed: item ${item.id} has no matching payload/source binding.`); + if (['applied', 'staged'].includes(prior.status)) return { replayed: true, digest }; + if (['applying', 'failed'].includes(prior.status)) throw new Error(`Batch precondition failed: item ${item.id} has an uncertain or in-flight outcome; inspect the backend before using a fresh journal.`); + } + const claim = { id: item.id, op: item.op, status: 'applying', digest, token: randomUUID() }; + writeBatchJournal(file, claim); + return claim; + }, { label: 'batch journal' }); +} + +export function finishBatchItem(file, claim, outcome) { + if (!file || !claim) return; + withLockSync(file, () => { + const current = readBatchJournal(file).get(outcome.id); + if (current?.token !== claim.token || current?.status !== 'applying') throw new Error('Batch precondition failed: applying claim changed.'); + writeBatchJournal(file, { ...outcome, digest: claim.digest, token: claim.token }); + }, { label: 'batch journal' }); } export function withTimeout(promise, timeoutMs, label = 'operation') { diff --git a/packages/cli/test/batch-resume.test.js b/packages/cli/test/batch-resume.test.js new file mode 100644 index 0000000..a7ddd98 --- /dev/null +++ b/packages/cli/test/batch-resume.test.js @@ -0,0 +1,59 @@ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { spawnSync } from 'node:child_process'; +import { pathToFileURL, fileURLToPath } from 'node:url'; +import { claimIdempotencyKey, recordIdempotencyKey } from '../idempotency.js'; +import { claimBatchItem, finishBatchItem } from '../reliability.js'; + +const cli = fileURLToPath(new URL('../bin/ats.js', import.meta.url)); +function temp(t) { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'ats-replay-')); + t.after(() => fs.rmSync(dir, { recursive: true, force: true })); return dir; +} + +test('create claims and batch claims refuse concurrent, changed and uncertain replay', (t) => { + const configDir = temp(t); + const request = { title: 'Ship', projectId: 'p' }; + const opts = { configDir, scope: 'source-a' }; + const claim = claimIdempotencyKey('same', request, opts); + assert.throws(() => claimIdempotencyKey('same', request, opts), /in flight or uncertain/); + recordIdempotencyKey('same', { taskId: 't', projectId: 'p' }, { ...opts, token: claim.entry.token }); + assert.equal(claimIdempotencyKey('same', request, opts).replayed, true); + assert.throws(() => claimIdempotencyKey('same', request, { ...opts, scope: 'source-b' }), /binding/); + const file = path.join(configDir, 'journal.jsonl'); + const item = { id: 'a', op: 'create', title: 'Ship' }; + const batch = claimBatchItem(file, item, 'source-a'); + assert.throws(() => claimBatchItem(file, item, 'source-a'), /in-flight/); + finishBatchItem(file, batch, { id: 'a', status: 'failed' }); + assert.throws(() => claimBatchItem(file, item, 'source-a'), /uncertain/); + assert.throws(() => claimBatchItem(file, { ...item, title: 'Changed' }, 'source-a'), /binding/); +}); + +test('CLI batch resume skips only unchanged operations; dry run leaves journal bytes untouched', (t) => { + const dir = temp(t); + const adapter = path.join(dir, 'adapter.mjs'); + const counter = path.join(dir, 'writes.jsonl'); + fs.writeFileSync(adapter, `import fs from 'node:fs';export default { + listProjects:async()=>[],listTasksInProject:async()=>[],getTask:async()=>null, + createTask:async(input)=>{fs.appendFileSync(${JSON.stringify(counter)},JSON.stringify(input)+'\\n');return {id:'t',...input};}, + updateTask:async()=>null,urlFor:()=>'',authStatus:async()=>({authenticated:true}),authLogin:async()=>({}) + };`); + const input = path.join(dir, 'batch.json'); const journal = path.join(dir, 'journal.jsonl'); + const run = (...flags) => spawnSync(process.execPath, [cli, 'batch', input, '--journal', journal, '--json', ...flags], { + encoding: 'utf8', env: { ...process.env, ATS_ADAPTER: pathToFileURL(adapter).href, XDG_CONFIG_HOME: path.join(dir, 'xdg'), ATS_ACTION_LOG: path.join(dir, 'actions.jsonl'), ATS_USAGE_DISABLE: '1' }, + }); + fs.writeFileSync(input, JSON.stringify([{ id: 'a', op: 'create', title: 'First', projectId: 'p' }])); + assert.equal(run().status, 0); + assert.equal(JSON.parse(run().stdout).summary.skipped, 1); + const bytes = fs.readFileSync(journal, 'utf8'); + assert.equal(run('--dry-run').status, 0); + assert.equal(fs.readFileSync(journal, 'utf8'), bytes); + fs.writeFileSync(input, JSON.stringify([{ id: 'a', op: 'create', title: 'Changed', projectId: 'p' }])); + const changed = run(); + assert.equal(changed.status, 5); + assert.match(changed.stdout, /binding/); + assert.equal(fs.readFileSync(counter, 'utf8').trim().split('\n').length, 1); +}); diff --git a/packages/cli/test/create-idempotent.test.js b/packages/cli/test/create-idempotent.test.js index 76a05a9..22d3961 100644 --- a/packages/cli/test/create-idempotent.test.js +++ b/packages/cli/test/create-idempotent.test.js @@ -122,3 +122,46 @@ test('keys age out and are pruned on write', () => { const store = JSON.parse(fs.readFileSync(path.join(configDir, 'idempotency-keys.json'), 'utf8')); assert.deepEqual(Object.keys(store.keys), ['fresh']); }); + + +test('a create key refuses a different payload and a corrupt store before writing', () => { + const before = creates(); + run('create', 'p1', 'Bound task', '--idempotency-key', 'bound'); + const changed = runProcess('create', 'p1', 'Different task', '--idempotency-key', 'bound'); + assert.equal(changed.status, 3); + assert.match(changed.stderr, /request\/source binding/); + assert.equal(creates(), before + 1); + const file = path.join(tempDir, 'xdg', 'ats', 'idempotency-keys.json'); + const saved = fs.readFileSync(file); + try { + fs.writeFileSync(file, '{'); + const bad = runProcess('create', 'p1', 'Corrupt', '--idempotency-key', 'new-key'); + assert.notEqual(bad.status, 0); + assert.equal(creates(), before + 1); + } finally { fs.writeFileSync(file, saved); } +}); + +test('a reviewed create keeps its binding through apply and subsequent replay', () => { + const previous = process.env.ATS_REVIEW_ALL; + process.env.ATS_REVIEW_ALL = '1'; + try { + const before = creates(); + const staged = run('create', 'p1', 'Reviewed task', '--idempotency-key', 'reviewed'); + assert.equal(staged.staged, true); + assert.equal(creates(), before); + run('review', 'approve', staged.reviewId); + const applied = run('review', 'apply', staged.reviewId); + assert.equal(applied.applied[0].ok, true); + assert.equal(creates(), before + 1); + const replay = run('create', 'p1', 'Reviewed task', '--idempotency-key', 'reviewed'); + assert.equal(replay.idempotent, true); + assert.equal(replay.task.id.startsWith('new-'), true); + assert.equal(creates(), before + 1); + const changed = runProcess('create', 'p1', 'Changed reviewed task', '--idempotency-key', 'reviewed'); + assert.equal(changed.status, 3); + assert.equal(creates(), before + 1); + } finally { + if (previous === undefined) delete process.env.ATS_REVIEW_ALL; + else process.env.ATS_REVIEW_ALL = previous; + } +}); diff --git a/packages/core/index.js b/packages/core/index.js index c155497..5198f7e 100644 --- a/packages/core/index.js +++ b/packages/core/index.js @@ -20,6 +20,8 @@ export { loadFacts, proposeFact, proposeRetract, + proposeConfirm, + staleFacts, ratifyFactItem, listKgFacts, askFacts, diff --git a/packages/core/kg-store.js b/packages/core/kg-store.js new file mode 100644 index 0000000..40678cb --- /dev/null +++ b/packages/core/kg-store.js @@ -0,0 +1,1007 @@ +/** + * `ats kg` — a facts layer beside the task layer. + * + * Agents accumulate durable, plain-language knowledge ("client X prefers Y", + * "service Z is deprecated") that outlives any single task. This module + * stores such knowledge as subject–predicate–object facts with temporal + * validity and provenance, embedded and serverless: an append-only JSONL + * event log under the config dir — no graph server, no daemon, nothing to + * operate. `ats state export` carries it like any other state file. + * + * Three deliberate properties, learned the hard way from running production + * knowledge graphs: + * + * 1. SINGLE WRITER. Nothing writes the fact store except `ratify`. + * Agents PROPOSE facts (and retractions); proposals stage in the same + * review queue as gated task writes (kind `kg.fact`) and reach the + * store only after a human approves. Every fact carries who proposed + * it, who ratified it, and from what source. + * 2. FACTS ARE EVENTS, NOT ROWS. The store is an append-only log of + * add/retract events, folded on read. A retraction closes a fact's + * validity interval (tInvalid) instead of deleting it — "what did we + * believe in June" stays answerable. + * 3. ZERO-LLM READS. `ask` is deterministic lexical scoring over the + * folded facts — fast, cheap, reproducible. Semantic retrieval can sit + * on top later; it is not required to get value. + * + * Domains partition the store like group ids ("sales", "infra"), so a + * question can be scoped to one domain or span all of them. For teams on an + * embedded graph database, `exportFactsCypher()` emits a script that loads + * the graph into LadybugDB/Kùzu-style engines (node table Entity, rel table + * FACT). + */ + +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { randomUUID } from 'node:crypto'; +import { stableDigest } from './reliability-snapshot.js'; +import { withLockSync } from './fs-lock.js'; +import { stageReviewItem, listReviewItems } from './review-queue.js'; + +export function kgFactsPath() { + const configBase = process.env.XDG_CONFIG_HOME || path.join(os.homedir(), '.config'); + return process.env.ATS_KG_FACTS || path.join(configBase, 'ats', 'kg-facts.jsonl'); +} + +// One lock, one write: events that belong together (close the old fact, add +// its replacement) land in a single append, so a crash cannot leave half a +// supersession in the log. +function appendEventsUnlocked(events, factsPath) { + fs.mkdirSync(path.dirname(factsPath), { recursive: true, mode: 0o700 }); + fs.appendFileSync(factsPath, events.map((event) => JSON.stringify(event) + '\n').join(''), { mode: 0o600 }); + fs.chmodSync(factsPath, 0o600); +} + +/** + * Fold the event log into current fact state. Closed facts stay — a + * `retract` closes validity with status `retracted`, a `supersede` closes it + * with status `superseded` and points at the replacement. A fact is closed + * once: a second closing event for the same fact is recorded in its history + * as ignored and never moves tInvalid. `history` maps fact id → the events + * that touched it, in log order. + */ +export function loadFacts({ factsPath = kgFactsPath() } = {}) { + if (!fs.existsSync(factsPath)) return { facts: [], events: 0, path: factsPath, history: new Map() }; + const lines = fs.readFileSync(factsPath, 'utf8').split('\n').filter(Boolean); + const byId = new Map(); + const history = new Map(); + const proposals = new Map(); + const note = (id, entry) => { + if (!history.has(id)) history.set(id, []); + history.get(id).push(entry); + }; + let events = 0; + for (const [index, line] of lines.entries()) { + let event; + try { + event = JSON.parse(line); + } catch (error) { + throw new Error(`Malformed kg fact log at line ${index + 1}: ${factsPath}`, { cause: error }); + } + events += 1; + const proposalId = event.proposalId || event.fact?.provenance?.proposalId; + if (proposalId) proposals.set(proposalId, event); + if (event.op === 'add' && event.fact?.id) { + byId.set(event.fact.id, { ...event.fact, status: 'active', tInvalid: null }); + note(event.fact.id, { + op: 'add', + at: event.at || event.fact.tValid || null, + by: event.fact.provenance?.ratifiedBy || null, + ...(event.fact.supersedes ? { supersedes: event.fact.supersedes } : {}), + }); + } else if (event.op === 'confirm' && byId.has(event.factId)) { + const fact = byId.get(event.factId); + const entry = { op: 'confirm', at: event.at, by: event.by, source: event.source }; + if (fact.status !== 'active') { note(event.factId, { ...entry, ignored: `already ${fact.status}` }); continue; } + fact.lastConfirmedAt = event.at; + fact.confirmation = { by: event.by, source: event.source, proposalId: event.proposalId }; + note(event.factId, entry); + } else if ((event.op === 'retract' || event.op === 'supersede') && byId.has(event.factId)) { + const fact = byId.get(event.factId); + const entry = { + op: event.op, + at: event.at || null, + by: event.by || null, + ...(event.reason ? { reason: event.reason } : {}), + ...(event.byFact ? { byFact: event.byFact } : {}), + }; + if (fact.status !== 'active') { + note(event.factId, { ...entry, ignored: `already ${fact.status}` }); + continue; + } + fact.status = event.op === 'retract' ? 'retracted' : 'superseded'; + fact.tInvalid = event.at || null; + if (event.reason) fact.retractReason = event.reason; + if (event.by) fact.retractedBy = event.by; + if (event.op === 'supersede') fact.supersededBy = event.byFact || null; + note(event.factId, entry); + } + } + return { facts: [...byId.values()], events, path: factsPath, history, proposals }; +} + +function requireText(value, label) { + if (!value || typeof value !== 'string' || !value.trim()) throw new Error(`kg: ${label} is required.`); + return value.trim(); +} + +/** + * Normalize an entity, predicate, or domain for comparison: case, runs of + * whitespace, and trailing punctuation do not make two mentions different. + */ +export function normalizeTerm(value) { + return String(value ?? '').toLowerCase().replace(/\s+/g, ' ').trim().replace(/[.,;:!?]+$/, ''); +} + +const SEP = '\u0000'; +const tripleKey = (f) => [f.domain || 'default', f.subject, f.predicate, f.object].map(normalizeTerm).join(SEP); +const pairKey = (f) => [f.domain || 'default', f.subject, f.predicate].map(normalizeTerm).join(SEP); +const factSummary = (f) => ({ id: f.id, subject: f.subject, predicate: f.predicate, object: f.object, domain: f.domain, tValid: f.tValid }); + +/** Thrown by `proposeFact` when the gate refuses; `code` is the verdict and `gate` the full report. */ +export class KgGateError extends Error { + constructor(gate) { + super(gate.message); + this.name = 'KgGateError'; + this.code = gate.verdict; + this.gate = gate; + } +} + +/** + * The proposal gate — pure, run before anything is staged. A proposal is + * checked against the folded store and the review queue and gets one verdict: + * + * duplicate — the same triple is already an active fact, or is already + * pending/approved for review (in-batch duplicates included) + * rejected — a human rejected this exact triple before; re-proposing it + * needs `acknowledgeRejected` naming that decision + * contradiction — active facts already state what this subject+predicate is, + * with another object; pass `supersedes` (replace one) or + * `additive` (the predicate holds several values) + * clear — stage it (conflicts, if any, are reported for the reviewer) + * + * Matching is normalized (case, whitespace, trailing punctuation) and scoped + * to the domain: the same triple in another domain is another graph. + */ +export function checkFactProposal(proposal, { facts = [], reviewItems = [], supersedes = null, additive = false, acknowledgeRejected = null } = {}) { + const key = tripleKey(proposal); + const pair = pairKey(proposal); + const shortId = (id) => String(id).slice(0, 8); + + const activeTwin = facts.find((f) => f.status === 'active' && tripleKey(f) === key); + if (activeTwin) { + return { + verdict: 'duplicate', + message: `kg: already an active fact (${shortId(activeTwin.id)}): ${activeTwin.subject} ${activeTwin.predicate} ${activeTwin.object}.`, + fact: factSummary(activeTwin), + }; + } + const queuedTwin = reviewItems.find((i) => i.kind === 'kg.fact' && (i.status === 'pending' || i.status === 'approved') + && i.payload?.op === 'add' && tripleKey(i.payload) === key); + if (queuedTwin) { + return { + verdict: 'duplicate', + message: `kg: the same fact is already ${queuedTwin.status} for review (${shortId(queuedTwin.id)}, proposed by ${queuedTwin.stagedBy}).`, + proposal: { id: queuedTwin.id, status: queuedTwin.status, stagedBy: queuedTwin.stagedBy, stagedAt: queuedTwin.stagedAt }, + }; + } + const rejected = reviewItems + .filter((i) => i.kind === 'kg.fact' && i.status === 'rejected' && i.payload?.op === 'add' && tripleKey(i.payload) === key) + .sort((a, b) => String(b.decidedAt || '').localeCompare(String(a.decidedAt || ''))); + if (rejected.length) { + const r = rejected[0]; + const acknowledged = acknowledgeRejected && (r.id === acknowledgeRejected || r.id.startsWith(acknowledgeRejected)); + if (!acknowledged) { + const when = String(r.decidedAt || '').slice(0, 10); + return { + verdict: 'rejected', + message: `kg: this fact was rejected${when ? ` on ${when}` : ''} by ${r.decidedBy || 'a reviewer'}${r.decisionNote ? ` (${r.decisionNote})` : ''} — review item ${shortId(r.id)}. If something changed, re-propose with --acknowledge-rejected ${shortId(r.id)}.`, + rejected: { id: r.id, decidedAt: r.decidedAt || null, decidedBy: r.decidedBy || null, note: r.decisionNote || null }, + }; + } + } + const conflicts = facts.filter((f) => f.status === 'active' && pairKey(f) === pair).map(factSummary); + if (conflicts.length && !supersedes && !additive) { + return { + verdict: 'contradiction', + message: `kg: ${conflicts.length} active fact${conflicts.length === 1 ? '' : 's'} already state${conflicts.length === 1 ? 's' : ''} what "${proposal.subject} ${proposal.predicate}" is: ${conflicts.map((f) => `${f.object} (${shortId(f.id)})`).join(', ')}. Pass --supersedes to replace one, or --additive if the predicate holds several values.`, + conflicts, + }; + } + return { verdict: 'clear', ...(conflicts.length ? { conflicts } : {}), ...(supersedes ? { supersedes } : {}), ...(additive ? { additive: true } : {}) }; +} + +/** + * Stage a fact proposal for review. Nothing becomes queryable here — a human + * approves (`ats review approve`) and `ats kg ratify` writes the store. + * + * The gate runs first (see `checkFactProposal`); a refused proposal throws + * `KgGateError`. `supersedes` names an active fact the new one replaces at + * ratification; `additive` allows a second value for the same + * subject+predicate; `acknowledgeRejected` re-opens a triple a human declined. + */ +export function proposeFact({ subject, predicate, object, domain, source, confidence, taskRef, by, supersedes, additive, acknowledgeRejected, validAt, learnedAt } = {}, { queuePath, factsPath } = {}) { + const payload = { + op: 'add', + ...(validAt !== undefined ? { validAt: factTimestamp(validAt, 'validAt') } : {}), + ...(learnedAt !== undefined ? { learnedAt: factTimestamp(learnedAt, 'learnedAt') } : {}), + subject: requireText(subject, 'subject'), + predicate: requireText(predicate, 'predicate'), + object: requireText(object, 'object'), + domain: (domain || 'default').trim(), + source: source || null, + confidence: confidence || 'medium', + ...(taskRef ? { taskRef } : {}), + }; + const { facts } = loadFacts(factsPath ? { factsPath } : {}); + if (supersedes) { + const target = facts.find((f) => f.id === supersedes || f.id.startsWith(supersedes)); + if (!target) throw new Error(`kg: --supersedes names no fact ${supersedes}.`); + if (target.status !== 'active') throw new Error(`kg: fact ${target.id} is already ${target.status}; it cannot be superseded.`); + payload.supersedes = target.id; + } + if (additive) payload.additive = true; + const reviewItems = listReviewItems({ kind: 'kg.fact', ...(queuePath ? { queuePath } : {}) }); + const gate = checkFactProposal(payload, { facts, reviewItems, supersedes: payload.supersedes || null, additive: !!additive, acknowledgeRejected }); + if (gate.verdict !== 'clear') throw new KgGateError(gate); + if (acknowledgeRejected) payload.acknowledgedRejection = acknowledgeRejected; + return stageReviewItem({ kind: 'kg.fact', payload, by, note: 'fact proposal' }, queuePath ? { queuePath } : {}); +} + +function parseTaskRef(value) { + if (!value) return undefined; + if (typeof value === 'object') return value.taskId ? { projectId: value.projectId, taskId: value.taskId } : undefined; + const s = String(value); + const i = s.lastIndexOf('/'); + if (i <= 0 || i === s.length - 1) throw new Error(`task must be PROJECT/TASK, got "${s}".`); + return { projectId: s.slice(0, i), taskId: s.slice(i + 1) }; +} + +/** + * Batch proposals: one JSON object per line (`{subject, predicate, object, + * domain?, source?, confidence?, task?, supersedes?, additive?, + * acknowledgeRejected?}`), or already-parsed objects. Every line is handled on + * its own — a malformed or refused line is reported and the rest still stage — + * and every staged line goes through the gate, so a duplicate inside the batch + * is caught against the line that was staged just before it. `defaults` + * (domain, source, confidence) fill in what a line does not carry. + */ +export function proposeFactLines(lines, { by, defaults = {}, queuePath, factsPath } = {}) { + const results = []; + const counts = { staged: 0, duplicate: 0, refused: 0, invalid: 0 }; + const paths = { ...(queuePath ? { queuePath } : {}), ...(factsPath ? { factsPath } : {}) }; + for (const [index, raw] of [...lines].entries()) { + const line = index + 1; + if (typeof raw === 'string' && !raw.trim()) continue; + let input; + try { + input = typeof raw === 'string' ? JSON.parse(raw) : raw; + if (!input || typeof input !== 'object' || Array.isArray(input)) throw new Error('expected a JSON object'); + } catch (err) { + counts.invalid += 1; + results.push({ line, ok: false, invalid: true, error: `line ${line}: ${err.message}` }); + continue; + } + try { + const item = proposeFact({ + validAt: input.validAt ?? defaults.validAt, + learnedAt: input.learnedAt ?? defaults.learnedAt, + subject: input.subject, + predicate: input.predicate, + object: input.object, + domain: input.domain || defaults.domain, + source: input.source || defaults.source, + confidence: input.confidence || defaults.confidence, + taskRef: parseTaskRef(input.task ?? input.taskRef), + by, + supersedes: input.supersedes, + additive: !!input.additive, + acknowledgeRejected: input.acknowledgeRejected, + }, paths); + counts.staged += 1; + results.push({ line, ok: true, reviewId: item.id, ...(item.payload.supersedes ? { supersedes: item.payload.supersedes } : {}) }); + } catch (err) { + if (err instanceof KgGateError) { + if (err.code === 'duplicate') counts.duplicate += 1; + else counts.refused += 1; + results.push({ line, ok: err.code === 'duplicate', verdict: err.code, message: err.message, ...(err.gate.conflicts ? { conflicts: err.gate.conflicts } : {}), ...(err.gate.rejected ? { rejected: err.gate.rejected } : {}) }); + } else { + counts.invalid += 1; + results.push({ line, ok: false, invalid: true, error: `line ${line}: ${err.message}` }); + } + } + } + return { ...counts, lines: results.length, results }; +} + +/** Stage a retraction proposal — the red pen goes through review too. */ +export function proposeRetract({ factId, reason, by } = {}, { queuePath, factsPath } = {}) { + requireText(factId, 'factId'); + const { facts } = loadFacts(factsPath ? { factsPath } : {}); + const fact = facts.find((f) => f.id === factId || f.id.startsWith(factId)); + if (!fact) throw new Error(`kg: no fact ${factId}.`); + if (fact.status !== 'active') throw new Error(`kg: fact ${fact.id} is already ${fact.status}.`); + const payload = { op: 'retract', factId: fact.id, reason: reason || null, summary: `${fact.subject} ${fact.predicate} ${fact.object}` }; + return stageReviewItem({ kind: 'kg.fact', payload, by, note: 'fact retraction' }, queuePath ? { queuePath } : {}); +} + +/** + * Apply one APPROVED kg.fact review item to the store — the single writer. + * Returns the fact (add), the fact plus the id it superseded (add with + * `supersedes`), or the closed fact id (retract). Closing is checked against + * the store at this moment, not at proposal time: a fact that is already + * retracted or superseded is refused here, so two retractions staged before + * either is ratified cannot move the first closing time. + */ +/** Validate explicit source timestamps; a planned date cannot become a fact. */ +function factTimestamp(value, label, now = new Date()) { + if (typeof value !== 'string' || !/^\d{4}-\d{2}-\d{2}T.*(?:Z|[+-]\d{2}:\d{2})$/.test(value)) throw new Error(`kg: ${label} must be an ISO timestamp with a timezone.`); + const date = new Date(value); + const [year, month, day] = value.slice(0, 10).split('-').map(Number); + if (!Number.isFinite(date.getTime()) || new Date(Date.UTC(year, month - 1, day)).getUTCDate() !== day) throw new Error(`kg: invalid ${label}.`); + if (date.getTime() > new Date(now).getTime()) throw new Error(`kg: ${label} cannot be a future or planned date.`); + return date.toISOString(); +} + +/** Stage fresh evidence for an active fact; confirmation still needs approval. */ +export function proposeConfirm({ factId, source, by } = {}, { factsPath, queuePath } = {}) { + requireText(factId, 'factId'); + const facts = loadFacts(factsPath ? { factsPath } : {}).facts; + const matches = facts.filter((fact) => fact.id === factId || fact.id.startsWith(factId)); + if (matches.length !== 1) throw new Error(`kg: fact ${factId} is missing or ambiguous.`); + const fact = matches[0]; + if (fact.status !== 'active') throw new Error(`kg: fact ${fact.id} is already ${fact.status}.`); + return stageReviewItem({ kind: 'kg.fact', by, note: 'fact confirmation', payload: { + op: 'confirm', factId: fact.id, source: requireText(source, 'source'), domain: fact.domain, + summary: `${fact.subject} ${fact.predicate} ${fact.object}`, + } }, queuePath ? { queuePath } : {}); +} + +/** Read-only freshness view. Age is a review signal, never an automatic retraction. */ +export function staleFacts({ days = 60, domain, limit = 50, now = Date.now(), factsPath } = {}) { + if (!Number.isFinite(days) || days < 0 || !Number.isInteger(limit) || limit < 1) throw new Error('kg: days must be nonnegative and limit must be a positive integer.'); + const active = listKgFacts({ domain, ...(factsPath ? { factsPath } : {}) }); + const stale = active.map((fact) => { + const lastConfirmedAt = fact.lastConfirmedAt || fact.tLearned || fact.provenance?.ratifiedAt || fact.tValid; + const epoch = Date.parse(lastConfirmedAt); + const ageDays = Number.isFinite(epoch) ? (Number(now) - epoch) / 86400000 : null; + return { ...fact, ageDays, freshness: ageDays === null ? 'unknown' : 'stale', lastConfirmedAt: lastConfirmedAt || null }; + }).filter((fact) => fact.ageDays === null || fact.ageDays >= days) + .sort((a, b) => (b.ageDays ?? Infinity) - (a.ageDays ?? Infinity) || a.id.localeCompare(b.id)); + return { days, scanned: active.length, count: stale.length, truncated: stale.length > limit, facts: stale.slice(0, limit) }; +} + +export function ratifyFactItem(item, { factsPath = kgFactsPath(), now = new Date() } = {}) { + if (item?.kind !== 'kg.fact') throw new Error(`kg: review item ${item?.id} is not a kg.fact proposal.`); + if (!['approved', 'applying'].includes(item.status)) throw new Error(`kg: review item ${item.id} is ${item.status}, not approved.`); + if (!item.approvedDigest || item.approvedDigest !== stableDigest({ kind: item.kind, payload: item.payload })) throw new Error('kg: no matching payload approval; stage and approve it again.'); + const at = new Date(now).toISOString(); + const p = item.payload || {}; + return withLockSync(factsPath, () => { + const { facts, proposals = new Map() } = loadFacts({ factsPath }); + const prior = proposals.get(item.id); + if (prior) return prior.op === 'add' + ? { op: 'add', fact: facts.find((fact) => fact.id === prior.fact.id), replayed: true } + : { op: prior.op, factId: prior.factId, replayed: true }; + const mustBeOpen = (id, verb) => { + const target = facts.find((fact) => fact.id === id); + if (!target) throw new Error(`kg: cannot ${verb} ${id}: no such fact.`); + if (target.status !== 'active') throw new Error(`kg: cannot ${verb} ${target.id}: it is already ${target.status} since ${target.tInvalid}.`); + return target; + }; + if (p.op === 'retract' || p.op === 'confirm') { + const target = mustBeOpen(p.factId, p.op); + if (p.op === 'confirm') requireText(p.source, 'source'); + appendEventsUnlocked([{ op: p.op, factId: target.id, at, by: item.decidedBy || null, + proposalId: item.id, ...(p.op === 'confirm' ? { source: p.source } : { reason: p.reason || null }) }], factsPath); + return { op: p.op, factId: target.id }; + } + if (p.op !== 'add') throw new Error(`kg: invalid proposal operation ${p.op}.`); + const gate = checkFactProposal(p, { facts, supersedes: p.supersedes, additive: p.additive }); + if (gate.verdict === 'duplicate') return { op: 'add', fact: facts.find((fact) => fact.id === gate.fact.id), duplicate: true }; + if (gate.verdict !== 'clear') throw new KgGateError(gate); + const fact = { + id: randomUUID(), subject: requireText(p.subject, 'subject'), predicate: requireText(p.predicate, 'predicate'), object: requireText(p.object, 'object'), + domain: p.domain || 'default', + tValid: p.validAt ? factTimestamp(p.validAt, 'validAt', now) : at, + tLearned: p.learnedAt ? factTimestamp(p.learnedAt, 'learnedAt', now) : item.stagedAt && item.stagedAt <= at ? item.stagedAt : at, + confidence: p.confidence || 'medium', + ...(p.taskRef ? { taskRef: p.taskRef } : {}), ...(p.supersedes ? { supersedes: p.supersedes } : {}), + provenance: { proposedBy: item.stagedBy || null, source: p.source || null, proposalId: item.id, ratifiedBy: item.decidedBy || null, ratifiedAt: at }, + }; + if (p.supersedes) { + const old = mustBeOpen(p.supersedes, 'supersede'); + if (fact.tValid < old.tValid) throw new Error('kg: supersession cannot predate the fact it replaces.'); + appendEventsUnlocked([ + { op: 'supersede', factId: old.id, at: fact.tValid, recordedAt: at, by: item.decidedBy || null, byFact: fact.id, reason: p.reason || `superseded by ${fact.id}` }, + { op: 'add', at, fact }, + ], factsPath); + return { op: 'add', fact, superseded: old.id }; + } + appendEventsUnlocked([{ op: 'add', at, fact }], factsPath); + return { op: 'add', fact }; + }, { label: 'kg facts' }); +} + +/** + * An `asOf` instant: an ISO timestamp, or a bare date meaning the end of that + * day (so the day a fact was ratified counts as believing it). Null when unset. + */ +export function parseAsOf(value) { + if (value == null || value === '' || value === true) return null; + const s = String(value).trim(); + const bare = /^\d{4}-\d{2}-\d{2}$/.test(s); + const d = new Date(bare ? `${s}T23:59:59.999Z` : s); + if (Number.isNaN(d.getTime())) throw new Error(`kg: --as-of needs an ISO date or timestamp (got "${value}").`); + return d.toISOString(); +} + +/** Was this fact believed at `asOf`? Validity is [tValid, tInvalid); a fact without tValid is taken as always valid. */ +export function believedAt(fact, asOf) { + if (fact.tValid && fact.tValid > asOf) return false; + if (fact.tInvalid && fact.tInvalid <= asOf) return false; + return true; +} + +// Does a fact touch an entity, as subject or object? Exact after +// normalization first ("acme gmbh" is "Acme GmbH"); when nothing matches +// exactly, a substring match ("Acme" finds "Acme GmbH") — the way --subject +// already reads. +function touching(facts, entity) { + const c = normalizeTerm(entity); + if (!c) return facts; + const exact = facts.filter((f) => normalizeTerm(f.subject) === c || normalizeTerm(f.object) === c); + return exact.length ? exact : facts.filter((f) => normalizeTerm(f.subject).includes(c) || normalizeTerm(f.object).includes(c)); +} + +export function listKgFacts({ domain, subject, predicate, entity, status = 'active', asOf, factsPath } = {}) { + const at = parseAsOf(asOf); + const { facts } = loadFacts(factsPath ? { factsPath } : {}); + const selected = facts + .filter((f) => (at ? believedAt(f, at) : status === 'all' ? true : f.status === status)) + .filter((f) => !domain || f.domain === domain); + return (entity ? touching(selected, entity) : selected) + .filter((f) => !subject || f.subject.toLowerCase().includes(subject.toLowerCase())) + .filter((f) => !predicate || f.predicate.toLowerCase().includes(predicate.toLowerCase())); +} + +/** + * The entity view: every subject and object the store knows, with how much + * it knows about each — fact count, which side it appears on, domains, + * predicates — most-known first. Spellings that normalize alike are one + * entity, shown under the spelling ratified first. `query` filters by name. + */ +export function listEntities({ domain, query, status = 'active', limit = 50, factsPath } = {}) { + const facts = listKgFacts({ domain, status, ...(factsPath ? { factsPath } : {}) }) + .sort((a, b) => String(a.tValid || '').localeCompare(String(b.tValid || ''))); + const byKey = new Map(); + const touch = (name, fact, side) => { + const key = normalizeTerm(name); + if (!key) return; + const entry = byKey.get(key) || { name, facts: 0, asSubject: 0, asObject: 0, domains: {}, predicates: {}, lastValid: null }; + entry.facts += 1; + entry[side] += 1; + entry.domains[fact.domain] = (entry.domains[fact.domain] || 0) + 1; + entry.predicates[fact.predicate] = (entry.predicates[fact.predicate] || 0) + 1; + if (!entry.lastValid || String(fact.tValid || '') > entry.lastValid) entry.lastValid = fact.tValid || null; + byKey.set(key, entry); + }; + for (const f of facts) { + touch(f.subject, f, 'asSubject'); + touch(f.object, f, 'asObject'); + } + const q = normalizeTerm(query); + const entities = [...byKey.values()] + .filter((e) => !q || normalizeTerm(e.name).includes(q)) + .map((e) => ({ ...e, predicates: Object.entries(e.predicates).sort((a, b) => b[1] - a[1]).map(([p]) => p) })) + .sort((a, b) => b.facts - a.facts || a.name.localeCompare(b.name)); + return { count: entities.length, entities: entities.slice(0, limit) }; +} + +function tokenize(text) { + return [...new Set(String(text || '').toLowerCase().match(/[a-z0-9]+/g) || [])]; +} + +/** + * The events that touched one fact, plus the supersession chain around it: + * `chain.replaces` walks back through what this fact superseded, + * `chain.replacedBy` walks forward to what superseded it. + */ +export function factHistory(factIdOrPrefix, { factsPath } = {}) { + requireText(factIdOrPrefix, 'factId'); + const { facts, history } = loadFacts(factsPath ? { factsPath } : {}); + const fact = facts.find((f) => f.id === factIdOrPrefix || f.id.startsWith(factIdOrPrefix)); + if (!fact) throw new Error(`kg: no fact ${factIdOrPrefix}.`); + const byId = new Map(facts.map((f) => [f.id, f])); + const walk = (field) => { + const ids = []; + for (let cur = fact; cur?.[field] && byId.has(cur[field]) && !ids.includes(cur[field]); cur = byId.get(cur[field])) ids.push(cur[field]); + return ids; + }; + return { fact, events: history.get(fact.id) || [], chain: { replaces: walk('supersedes'), replacedBy: walk('supersededBy') } }; +} + +/** The facts a question may be answered from: by validity at `at`, else by status; anchored on `center` when given. */ +function selectFacts(facts, { domain, includeRetracted = false, at = null, center = null } = {}) { + const selected = facts + .filter((f) => (at ? believedAt(f, at) : includeRetracted || f.status === 'active')) + .filter((f) => !domain || f.domain === domain); + return center ? touching(selected, center) : selected; +} + +// How much of the question the top lexical hit carries decides the verdict: +// an agent reads it before treating one fact as the answer. +function lexicalConfidence(scored, tokens) { + if (!scored.length) return { verdict: 'none', reason: 'no fact shares a term with the question', coverage: 0 }; + const top = scored[0]; + const coverage = tokens.length ? top.hits / tokens.length : 0; + const ties = scored.filter((e) => e.score === top.score).length - 1; + let verdict; + let reason; + if (top.phrase || coverage >= 0.75) { + verdict = 'strong'; + reason = top.phrase ? 'the question names the fact' : 'the top fact carries most of the question'; + } else if (coverage >= 0.5) { + verdict = 'moderate'; + reason = 'the top fact carries half the question'; + } else { + verdict = 'weak'; + reason = 'the top fact shares one or two terms with the question — refine the question, scope it, or ask --semantic'; + } + return { verdict, reason, coverage: Math.round(coverage * 100) / 100, ...(ties ? { ties } : {}) }; +} + +/** + * Deterministic, zero-LLM question answering over the folded facts: + * token overlap weighted subject > object > predicate, exact-phrase bonus, + * newest-first tiebreak. Returns scored facts with full provenance and a + * `confidence` verdict (`strong` / `moderate` / `weak` / `none`) read from + * how much of the question the top fact carries. + * `asOf` answers from the validity intervals instead of the current status — + * what the store believed at that instant, closed facts included. `center` + * anchors the answer on one entity: only facts touching it are candidates. + */ +export function askFacts(question, { domain, limit = 8, includeRetracted = false, asOf, center, factsPath } = {}) { + const at = parseAsOf(asOf); + const tokens = tokenize(question); + const { facts, path: p } = loadFacts(factsPath ? { factsPath } : {}); + const scope = { ...(at ? { asOf: at } : {}), ...(center ? { center } : {}) }; + if (tokens.length === 0) { + return { question, count: 0, path: p, ...scope, confidence: lexicalConfidence([], tokens), facts: [] }; + } + const phrase = String(question || '').toLowerCase().trim(); + const scored = selectFacts(facts, { domain, includeRetracted, at, center }) + .map((f) => { + const s = f.subject.toLowerCase(); + const pr = f.predicate.toLowerCase(); + const o = f.object.toLowerCase(); + let score = 0; + let hits = 0; + for (const token of tokens) { + let hit = false; + if (s.includes(token)) { score += 3; hit = true; } + if (o.includes(token)) { score += 2; hit = true; } + if (pr.includes(token)) { score += 1; hit = true; } + if (hit) hits += 1; + } + if (score > 0) score /= tokens.length; + const phraseHit = !!phrase && (s.includes(phrase) || o.includes(phrase)); + if (phraseHit) score += 2; + return { fact: f, score, hits, phrase: phraseHit }; + }) + .filter((entry) => entry.score > 0) + .sort((a, b) => b.score - a.score || String(b.fact.tValid).localeCompare(String(a.fact.tValid))) + .slice(0, limit); + const confidence = lexicalConfidence(scored, tokens); + const out = scored.map(({ fact, score }) => ({ score: Math.round(score * 100) / 100, ...fact })); + return { question, count: out.length, path: p, ...scope, confidence, facts: out }; +} + +export function kgVectorsPath() { + const configBase = process.env.XDG_CONFIG_HOME || path.join(os.homedir(), '.config'); + return process.env.ATS_KG_VECTORS || path.join(configBase, 'ats', 'kg-vectors.json'); +} + +const factText = (f) => `${f.subject} ${f.predicate} ${f.object}`; + +function cosine(a, b) { + let dot = 0; + let na = 0; + let nb = 0; + const n = Math.min(a.length, b.length); + for (let i = 0; i < n; i += 1) { + dot += a[i] * b[i]; + na += a[i] * a[i]; + nb += b[i] * b[i]; + } + return na && nb ? dot / (Math.sqrt(na) * Math.sqrt(nb)) : 0; +} + +function readVectorCache(vectorsPath, cacheKey) { + try { + const parsed = JSON.parse(fs.readFileSync(vectorsPath, 'utf8')); + if (parsed && parsed.cacheKey === cacheKey && parsed.entries && typeof parsed.entries === 'object') return parsed; + } catch { + // no cache, or a cache for another embedder: start over + } + return { version: 1, cacheKey, entries: {} }; +} + +function writeVectorCache(vectorsPath, cache) { + withLockSync(vectorsPath, () => { + fs.mkdirSync(path.dirname(vectorsPath), { recursive: true, mode: 0o700 }); + fs.writeFileSync(vectorsPath, JSON.stringify(cache), { mode: 0o600 }); + }, { label: 'kg vectors' }); +} + +/** + * The opt-in semantic ask: the lexical branch fused (reciprocal rank fusion) + * with a dense branch over the same facts, embedded by the caller's + * `embed(texts)` — typically the active adapter's `embeddings`. Fact vectors + * are cached under `cacheKey` (one cache per embedder; a different key starts + * over), so a repeat ask embeds only the question and any fact that changed. + * A failing embedder does not fail the ask: the result degrades to the + * lexical branch and says so (`degraded`, `branches`), and `confidence` reads + * branch agreement on the top fact the way `find` does. + */ +export async function askFactsSemantic(question, { + embed, cacheKey = 'default', vectorsPath = kgVectorsPath(), domain, limit = 8, includeRetracted = false, asOf, center, factsPath, +} = {}) { + if (typeof embed !== 'function') throw new Error('kg: askFactsSemantic needs an embed(texts) function.'); + const at = parseAsOf(asOf); + const depth = Math.max(limit, 20); + const lexical = askFacts(question, { domain, limit: depth, includeRetracted, asOf, center, factsPath }); + const { facts, path: p } = loadFacts(factsPath ? { factsPath } : {}); + const candidates = selectFacts(facts, { domain, includeRetracted, at, center }); + const branches = [{ name: 'lexical', ok: true, count: lexical.count }]; + let dense = []; + let denseError = null; + try { + const cache = readVectorCache(vectorsPath, cacheKey); + const missing = candidates.filter((f) => !cache.entries[f.id] || cache.entries[f.id].text !== factText(f)); + const vectors = await embed([question, ...missing.map(factText)]); + if (!Array.isArray(vectors) || vectors.length !== missing.length + 1) { + throw new Error(`embed returned ${Array.isArray(vectors) ? vectors.length : 'no'} vectors for ${missing.length + 1} texts`); + } + const [questionVector, ...factVectors] = vectors; + missing.forEach((f, i) => { cache.entries[f.id] = { text: factText(f), vector: factVectors[i] }; }); + if (missing.length) writeVectorCache(vectorsPath, cache); + dense = candidates + .map((f) => ({ fact: f, similarity: cosine(questionVector, cache.entries[f.id].vector) })) + .filter((entry) => entry.similarity > 0) + .sort((a, b) => b.similarity - a.similarity) + .slice(0, depth); + branches.push({ name: 'dense', ok: true, count: dense.length, embedded: missing.length, cached: candidates.length - missing.length }); + } catch (err) { + denseError = err.message; + branches.push({ name: 'dense', ok: false, count: 0, error: err.message }); + } + + const K = 60; + const fused = new Map(); + const merge = (list, name, extra) => { + list.forEach((entry, rank) => { + const current = fused.get(entry.fact.id) || { fact: entry.fact, rrf: 0, sources: [] }; + current.rrf += 1 / (K + rank + 1); + current.sources.push(name); + Object.assign(current, extra(entry)); + fused.set(entry.fact.id, current); + }); + }; + merge(lexical.facts.map((f) => ({ fact: f })), 'lexical', (e) => ({ lexicalScore: e.fact.score })); + merge(dense, 'dense', (e) => ({ similarity: Math.round(e.similarity * 1000) / 1000 })); + const ranked = [...fused.values()] + .sort((a, b) => b.rrf - a.rrf || String(b.fact.tValid).localeCompare(String(a.fact.tValid))) + .slice(0, limit) + .map(({ fact, rrf, sources, lexicalScore, similarity }) => { + const { score: _lexical, ...rest } = fact; + return { + score: Math.round(rrf * 10000) / 10000, + sources, + ...(lexicalScore !== undefined ? { lexicalScore } : {}), + ...(similarity !== undefined ? { similarity } : {}), + ...rest, + }; + }); + + const topLexical = lexical.facts[0]?.id; + const topDense = dense[0]?.fact.id; + const branchesRun = denseError ? 1 : 2; + let confidence; + if (!ranked.length) { + confidence = { verdict: 'none', reason: denseError ? `no fact shares a term with the question, and the dense branch failed (${denseError})` : 'neither branch found a fact', branchesRun, topAgreement: 0 }; + } else if (!denseError && topLexical && topLexical === topDense) { + confidence = { verdict: 'strong', reason: 'both branches rank the same fact first', branchesRun, topAgreement: 2 }; + } else if (ranked[0].sources.length >= 2) { + confidence = { verdict: 'moderate', reason: 'both branches found the top fact, ranked differently', branchesRun, topAgreement: 2 }; + } else if (denseError) { + confidence = { + verdict: lexical.confidence.verdict === 'strong' ? 'moderate' : 'weak', + reason: `dense branch failed (${denseError}); lexical only: ${lexical.confidence.reason}`, + branchesRun, + topAgreement: 1, + }; + } else { + confidence = { verdict: 'weak', reason: `only the ${ranked[0].sources[0]} branch found the top fact`, branchesRun, topAgreement: 1 }; + } + return { + question, + mode: 'semantic', + count: ranked.length, + path: p, + ...(at ? { asOf: at } : {}), + ...(center ? { center } : {}), + degraded: !!denseError, + branches, + confidence, + facts: ranked, + }; +} + +const sameRef = (a, b) => { + if (a == null || b == null) return false; + const x = String(a); + const y = String(b); + return x === y || x.endsWith(`:${y}`) || y.endsWith(`:${x}`); +}; + +/** + * The facts that belong in a task's execution context: `linked` are active + * facts proposed from this task (`--task`), `related` are the best lexical + * matches for the task's title and intent. Each fact keeps its provenance and + * says how it got here (`via`). + */ +export function factsForTask({ projectId, taskId, query, domain, limit = 5, factsPath } = {}) { + const { facts } = loadFacts(factsPath ? { factsPath } : {}); + const active = facts.filter((f) => f.status === 'active' && (!domain || f.domain === domain)); + const linked = active + .filter((f) => f.taskRef && sameRef(f.taskRef.taskId, taskId) && (!projectId || !f.taskRef.projectId || sameRef(f.taskRef.projectId, projectId))) + .map((f) => ({ ...f, via: 'task-ref' })); + const linkedIds = new Set(linked.map((f) => f.id)); + let related = []; + if (query && String(query).trim()) { + related = askFacts(query, { domain, limit: limit + linked.length, ...(factsPath ? { factsPath } : {}) }).facts + .filter((f) => !linkedIds.has(f.id)) + .slice(0, limit) + .map((f) => ({ ...f, via: 'lexical' })); + } + return { linked, related, count: linked.length + related.length }; +} + +/** + * The reviewer's view of the queue: every fact proposal that is not in the + * store yet (pending, or approved and not ratified), with what ratifying it + * would do (`effect`) and how it reads against the store right now + * (`verdict`): `clear`, `duplicate` (its twin got ratified meanwhile), + * `contradiction` (an active fact arrived since it was staged), `rejected` + * (the same triple was declined since), or `stale` (the fact it retracts or + * supersedes is already closed). Grouped by domain so a reviewer sees what a + * `kg ratify --all` would promote into each graph. + */ +export function pendingFactProposals({ domain, factsPath, queuePath } = {}) { + const { facts } = loadFacts(factsPath ? { factsPath } : {}); + const items = listReviewItems({ kind: 'kg.fact', ...(queuePath ? { queuePath } : {}) }); + const open = items.filter((i) => i.status === 'pending' || i.status === 'approved'); + const shortId = (id) => String(id).slice(0, 8); + const pending = []; + for (const item of open) { + const p = item.payload || {}; + const base = { + id: item.id, + status: item.status, + op: p.op || 'add', + stagedBy: item.stagedBy || null, + stagedAt: item.stagedAt || null, + ...(item.decidedBy ? { approvedBy: item.decidedBy, approvedAt: item.decidedAt || null } : {}), + }; + if (p.op === 'retract' || p.op === 'confirm') { + const target = facts.find((f) => f.id === p.factId); + if (domain && target && target.domain !== domain) continue; + const stale = !target ? 'no such fact' : target.status !== 'active' ? `already ${target.status} since ${target.tInvalid}` : null; + pending.push({ + ...base, + domain: target?.domain || null, + effect: `${p.op} ${shortId(p.factId)}`, + target: target ? factSummary(target) : { id: p.factId, summary: p.summary || null }, + reason: p.reason || null, + ...(p.op === 'confirm' ? { source: p.source } : {}), + verdict: stale ? 'stale' : 'clear', + ...(stale ? { message: `kg: fact ${shortId(p.factId)} is ${stale}.` } : {}), + }); + continue; + } + const factDomain = p.domain || 'default'; + if (domain && factDomain !== domain) continue; + const entry = { + ...base, + domain: factDomain, + subject: p.subject, + predicate: p.predicate, + object: p.object, + source: p.source || null, + confidence: p.confidence || 'medium', + ...(p.taskRef ? { taskRef: p.taskRef } : {}), + ...(p.supersedes ? { supersedes: p.supersedes } : {}), + ...(p.additive ? { additive: true } : {}), + effect: p.supersedes ? `add, superseding ${shortId(p.supersedes)}` : p.additive ? 'add (additional value)' : 'add', + }; + if (p.supersedes) { + const target = facts.find((f) => f.id === p.supersedes); + const stale = !target ? 'no such fact' : target.status !== 'active' ? `already ${target.status} since ${target.tInvalid}` : null; + if (stale) { + pending.push({ ...entry, verdict: 'stale', message: `kg: fact ${shortId(p.supersedes)} is ${stale}; this proposal can no longer supersede it.` }); + continue; + } + } + const others = items.filter((i) => i.id !== item.id); + const gate = checkFactProposal(p, { facts, reviewItems: others, supersedes: p.supersedes || null, additive: !!p.additive, acknowledgeRejected: p.acknowledgedRejection || null }); + pending.push({ + ...entry, + verdict: gate.verdict, + ...(gate.message ? { message: gate.message } : {}), + ...(gate.conflicts ? { conflicts: gate.conflicts } : {}), + ...(gate.fact ? { duplicateOf: gate.fact } : {}), + ...(gate.rejected ? { rejected: gate.rejected } : {}), + }); + } + pending.sort((a, b) => String(a.domain).localeCompare(String(b.domain)) || String(a.stagedAt).localeCompare(String(b.stagedAt))); + const byDomain = {}; + for (const entry of pending) byDomain[entry.domain] = (byDomain[entry.domain] || 0) + 1; + return { count: pending.length, byDomain, pending }; +} + +export function kgStats({ factsPath, listReviewItems } = {}) { + const { facts, events, path: p } = loadFacts(factsPath ? { factsPath } : {}); + const domains = {}; + let active = 0; + let superseded = 0; + for (const f of facts) { + domains[f.domain] = (domains[f.domain] || 0) + 1; + if (f.status === 'active') active += 1; + if (f.status === 'superseded') superseded += 1; + } + const stats = { path: p, events, facts: facts.length, active, retracted: facts.length - active - superseded, superseded, domains }; + if (typeof listReviewItems === 'function') { + stats.pendingProposals = listReviewItems({ kind: 'kg.fact', status: 'pending' }).length; + stats.approvedUnratified = listReviewItems({ kind: 'kg.fact', status: 'approved' }).length; + } + return stats; +} + +function cypherEscape(value) { + return String(value ?? '').replace(/\\/g, '\\\\').replace(/'/g, "\\'"); +} + +// Every FACT relationship property in the Cypher export, in DDL order. The +// whole provenance record travels with the fact: who proposed it, who ratified +// it and when, the source, the task it came from, and — for a retracted fact — +// when its validity closed, by whom, and why. +const CYPHER_FACT_PROPS = [ + ['id', (f) => f.id], + ['predicate', (f) => f.predicate], + ['domain', (f) => f.domain], + ['status', (f) => f.status || 'active'], + ['tValid', (f) => f.tValid], + ['tLearned', (f) => f.tLearned], + ['lastConfirmedAt', (f) => f.lastConfirmedAt], + ['confirmationSource', (f) => f.confirmation?.source], + ['tInvalid', (f) => f.tInvalid], + ['confidence', (f) => f.confidence], + ['source', (f) => f.provenance?.source], + ['proposedBy', (f) => f.provenance?.proposedBy], + ['proposalId', (f) => f.provenance?.proposalId], + ['ratifiedBy', (f) => f.provenance?.ratifiedBy], + ['ratifiedAt', (f) => f.provenance?.ratifiedAt], + ['taskRef', (f) => f.taskRef], + ['retractedBy', (f) => f.retractedBy], + ['retractReason', (f) => f.retractReason], + ['supersedes', (f) => f.supersedes], + ['supersededBy', (f) => f.supersededBy], +]; + +export const CYPHER_DIALECTS = ['ladybug', 'kuzu', 'neo4j', 'falkordb']; + +/** + * Emit a Cypher script that loads the facts into a graph database. Active + * facts by default; `includeRetracted` adds the closed ones (retracted and + * superseded) with their closed validity, so the script is a complete record + * of what the store believed and why. + * + * Dialects: `ladybug` (default; `kuzu` is the same) speaks the embedded + * engine's typed DDL — node table Entity, rel table FACT — then MERGEs + * entities and CREATEs relationships. `neo4j` and `falkordb` speak + * openCypher with no table DDL: an index on Entity.name, entities MERGEd by + * name, and every fact MERGEd by its id with its properties SET — so the + * same script can be run again after facts close or arrive. + */ +export function exportFactsCypher({ domain, factsPath, includeRetracted = false, dialect = 'ladybug' } = {}) { + const d = String(dialect || 'ladybug').toLowerCase(); + if (!CYPHER_DIALECTS.includes(d)) throw new Error(`kg: dialect must be one of ${CYPHER_DIALECTS.join(', ')} (got "${dialect}").`); + const facts = listKgFacts({ domain, status: includeRetracted ? 'all' : 'active', ...(factsPath ? { factsPath } : {}) }); + const scope = `${includeRetracted ? 'all facts (active and retracted)' : 'active facts'} as a property graph, full provenance on every FACT.`; + const entities = new Set(); + for (const f of facts) { + entities.add(f.subject); + entities.add(f.object); + } + const entityLines = [...entities].sort().map((name) => `MERGE (:Entity {name: '${cypherEscape(name)}'});`); + const endpoints = (f) => `MATCH (a:Entity {name: '${cypherEscape(f.subject)}'}), (b:Entity {name: '${cypherEscape(f.object)}'}) `; + + if (d === 'ladybug' || d === 'kuzu') { + const ddlProps = CYPHER_FACT_PROPS.map(([name]) => `${name} STRING`).join(', '); + const lines = [ + `// ats kg export — ${scope}`, + '// Load with an embedded Cypher engine (LadybugDB / Kùzu): run the DDL once, then the data.', + "CREATE NODE TABLE IF NOT EXISTS Entity(name STRING, PRIMARY KEY(name));", + `CREATE REL TABLE IF NOT EXISTS FACT(FROM Entity TO Entity, ${ddlProps});`, + '', + ...entityLines, + '', + ]; + for (const f of facts) { + const props = CYPHER_FACT_PROPS.map(([name, read]) => `${name}: '${cypherEscape(read(f) ?? '')}'`).join(', '); + lines.push(`${endpoints(f)}CREATE (a)-[:FACT {${props}}]->(b);`); + } + return lines.join('\n') + '\n'; + } + + const engine = d === 'neo4j' ? 'Neo4j' : 'FalkorDB'; + const lines = [ + `// ats kg export — ${scope}`, + `// openCypher for ${engine}: re-runnable — entities MERGE by name, facts MERGE by id and SET their properties.`, + ...(d === 'neo4j' + ? ['CREATE INDEX entity_name IF NOT EXISTS FOR (e:Entity) ON (e.name);'] + : ['// FalkorDB refuses to create an index that exists — drop the next line on a re-run.', 'CREATE INDEX FOR (e:Entity) ON (e.name);']), + '', + ...entityLines, + '', + ]; + for (const f of facts) { + const sets = CYPHER_FACT_PROPS.filter(([name]) => name !== 'id').map(([name, read]) => `r.${name} = '${cypherEscape(read(f) ?? '')}'`).join(', '); + lines.push(`${endpoints(f)}MERGE (a)-[r:FACT {id: '${cypherEscape(f.id)}'}]->(b) SET ${sets};`); + } + return lines.join('\n') + '\n'; +} + +/** + * Emit the facts as Graphiti episodes, one JSON object per line, ready for an + * `add_episode` / `add_episode_bulk` ingest: `name`, `content` (the fact as a + * sentence), `source: "text"`, `source_description` (the provenance record), + * `reference_time` (tValid), plus `group_id` (the domain), `uuid` (the fact + * id) and the validity fields for a pipeline that wants them. A closed fact + * says so in its sentence, so the graph engine learns the retraction too. + */ +export function exportFactsGraphiti({ domain, factsPath, includeRetracted = false } = {}) { + const facts = listKgFacts({ domain, status: includeRetracted ? 'all' : 'active', ...(factsPath ? { factsPath } : {}) }); + const lines = facts.map((f) => { + let content = `${f.subject} ${f.predicate} ${f.object}.`; + if (f.status === 'retracted') content += ` (retracted${f.tInvalid ? ` ${f.tInvalid.slice(0, 10)}` : ''}${f.retractReason ? `: ${f.retractReason}` : ''})`; + if (f.status === 'superseded') content += ` (superseded${f.tInvalid ? ` ${f.tInvalid.slice(0, 10)}` : ''}${f.supersededBy ? ` by fact ${f.supersededBy}` : ''})`; + const prov = f.provenance || {}; + const description = [ + `ats kg fact ${f.id}`, + `domain ${f.domain}`, + prov.proposedBy ? `proposed by ${prov.proposedBy}${prov.source ? ` from ${prov.source}` : ''}` : (prov.source ? `source ${prov.source}` : null), + prov.ratifiedBy ? `ratified by ${prov.ratifiedBy}${prov.ratifiedAt ? ` ${prov.ratifiedAt}` : ''}` : null, + f.confidence ? `confidence ${f.confidence}` : null, + f.taskRef ? `task ${f.taskRef.projectId}/${f.taskRef.taskId}` : null, + ].filter(Boolean).join('; '); + return JSON.stringify({ + name: `ats-kg ${f.id.slice(0, 8)}`, + content, + source: 'text', + source_description: description, + reference_time: f.tValid || prov.ratifiedAt || null, + group_id: f.domain, + uuid: f.id, + status: f.status || 'active', + valid_at: f.tValid || null, + learned_at: f.tLearned || prov.ratifiedAt || null, + last_confirmed_at: f.lastConfirmedAt || null, + invalid_at: f.tInvalid || null, + }); + }); + return lines.length ? lines.join('\n') + '\n' : ''; +} diff --git a/packages/core/kg.js b/packages/core/kg.js index 65f240a..f8adf0f 100644 Binary files a/packages/core/kg.js and b/packages/core/kg.js differ diff --git a/packages/core/package.json b/packages/core/package.json index 27ff37f..8d96876 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-core", - "version": "0.15.0", + "version": "0.16.0", "description": "Adapter-agnostic core for Agentic Task System — retrieval (RRF parallel fan-out), corpus cache, usage logging, conformance kit, adapter interface", "type": "module", "main": "index.js", diff --git a/packages/core/retry.d.ts b/packages/core/retry.d.ts index 1dca5fa..dcbb20b 100644 --- a/packages/core/retry.d.ts +++ b/packages/core/retry.d.ts @@ -6,6 +6,9 @@ export interface RetryPolicy { } export interface RetryOptions extends Partial { + signal?: AbortSignal; + timeoutMs?: number; + retrySafe?: boolean; label?: string; sleep?: (ms: number) => Promise; onRetry?: (info: { attempt: number; waitMs: number; reason: string; label?: string }) => void; @@ -23,4 +26,4 @@ export function withRetry(attempt: (n: number) => Promise | T, opts?: Retr export function retryingFetch( fetchFn?: (url: string | URL, init?: RequestInit) => Promise, opts?: RetryOptions -): (url: string | URL, init?: RequestInit) => Promise; +): (url: string | URL, init?: RequestInit & { retrySafe?: boolean }) => Promise; diff --git a/packages/core/retry.js b/packages/core/retry.js index 84c5746..aa70f37 100644 --- a/packages/core/retry.js +++ b/packages/core/retry.js @@ -1,21 +1,16 @@ /** - * One retry policy for every adapter's HTTP path. + * Shared adapter HTTP retry policy. Reads retry transient network/gateway + * failures and rate limits. Mutations retry explicit rate-limit rejection only; + * a read-only POST can opt in with retrySafe. The request deadline covers + * attempts, backoff and response-body reads, and caller cancellation stops both + * requests and waits. Server waits above 60s return the original response. * - * Transient upstream conditions — 429, 502/503/504, a 500 whose body names a - * query/rate limit (TickTick's `exceed_query_limit`), a 403 that carries a - * Retry-After or an exhausted rate-limit window (GitHub secondary limits), and - * network-level failures (reset, timeout, DNS) — are retried with a jittered - * exponential backoff that honors `Retry-After` and `x-ratelimit-reset`. - * Everything else returns or throws on the first attempt: a 4xx that is not a - * rate limit is the caller's problem, never retried. - * - * Knobs (env): ATS_HTTP_RETRIES (default 3, 0 disables), ATS_HTTP_RETRY_BASE_MS - * (500), ATS_HTTP_RETRY_MAX_MS (8000). A single wait never exceeds 60s even when - * the upstream asks for more — an agent call that would block longer than that - * should fail loudly instead. + * Env: ATS_HTTP_RETRIES (3), ATS_HTTP_RETRY_BASE_MS (500), + * ATS_HTTP_RETRY_MAX_MS (8000), ATS_HTTP_TIMEOUT_MS (30000, positive). */ const RATE_LIMIT_BODY = /exceed_query_limit|rate.?limit|too many requests|quota exceeded|try again later|temporarily unavailable/i; +const EXPLICIT_RATE_LIMIT_BODY = /exceed_query_limit|rate.?limit|too many requests|quota exceeded/i; const TRANSIENT_NET = /ECONNRESET|ECONNREFUSED|ETIMEDOUT|EAI_AGAIN|ENOTFOUND|EPIPE|UND_ERR_SOCKET|UND_ERR_CONNECT_TIMEOUT|fetch failed|socket hang up|network/i; const MAX_WAIT_MS = 60_000; @@ -91,7 +86,31 @@ export function backoffMs(attempt, policy = retryPolicy()) { return Math.round(raw * (0.5 + Math.random() * 0.5)); } -const defaultSleep = (ms) => new Promise((r) => setTimeout(r, ms)); +function abortReason(signal) { + return signal?.reason || Object.assign(new Error('Request cancelled'), { name: 'AbortError' }); +} + +function abortable(promise, signal) { + if (!signal) return Promise.resolve(promise); + if (signal.aborted) { + Promise.resolve(promise).catch(() => {}); + return Promise.reject(abortReason(signal)); + } + return new Promise((resolve, reject) => { + const onAbort = () => { cleanup(); reject(abortReason(signal)); }; + const cleanup = () => signal.removeEventListener('abort', onAbort); + signal.addEventListener('abort', onAbort, { once: true }); + Promise.resolve(promise).then((value) => { cleanup(); resolve(value); }, (error) => { cleanup(); reject(error); }); + }); +} + +const defaultSleep = (ms, signal) => new Promise((resolve, reject) => { + if (signal?.aborted) { reject(abortReason(signal)); return; } + const cleanup = () => signal?.removeEventListener('abort', onAbort); + const timer = setTimeout(() => { cleanup(); resolve(); }, ms); + const onAbort = () => { clearTimeout(timer); cleanup(); reject(abortReason(signal)); }; + signal?.addEventListener('abort', onAbort, { once: true }); +}); async function peekBody(res) { try { @@ -116,10 +135,12 @@ export async function withRetry(attempt, opts = {}) { const sleep = opts.sleep || defaultSleep; const label = opts.label ? `${opts.label}: ` : ''; for (let n = 0; ; n++) { + if (opts.signal?.aborted) throw abortReason(opts.signal); let res; try { - res = await attempt(n); + res = await abortable(attempt(n), opts.signal); } catch (err) { + if (opts.signal?.aborted || err?.name === 'AbortError' || err?.name === 'TimeoutError') throw err; const retryable = opts.isRetryableError ? opts.isRetryableError(err) : isTransientError(err); if (!retryable || n >= policy.retries) { if (retryable && n > 0) err.message = `${err.message} (after ${n} retr${n === 1 ? 'y' : 'ies'})`; @@ -127,19 +148,21 @@ export async function withRetry(attempt, opts = {}) { } const wait = backoffMs(n, policy); opts.onRetry?.({ attempt: n + 1, waitMs: wait, reason: err.message, label: opts.label }); - await sleep(wait); + await abortable(sleep(wait, opts.signal), opts.signal); continue; } - const body = res && typeof res === 'object' && res.status === 500 ? await peekBody(res) : ''; + const body = res && typeof res === 'object' && res.status === 500 ? await abortable(peekBody(res), opts.signal) : ''; const transient = opts.isRetryableResponse ? opts.isRetryableResponse(res, body) : isTransientResponse(res, body); if (!transient || n >= policy.retries) { if (transient && res && typeof res === 'object') res.__retries = n; return res; } const asked = upstreamWaitMs(res); - const wait = Math.min(MAX_WAIT_MS, asked ?? backoffMs(n, policy)); + // Do not retry sooner than a server's explicit reset time. + if (asked !== null && asked > MAX_WAIT_MS) return res; + const wait = asked ?? backoffMs(n, policy); opts.onRetry?.({ attempt: n + 1, waitMs: wait, reason: `${label}HTTP ${res.status}`, label: opts.label }); - await sleep(wait); + await abortable(sleep(wait, opts.signal), opts.signal); } } @@ -148,5 +171,43 @@ export async function withRetry(attempt, opts = {}) { * Drop-in: `const fetch = retryingFetch(globalThis.fetch, { label: 'Notion' })`. */ export function retryingFetch(fetchFn = globalThis.fetch, opts = {}) { - return (url, init) => withRetry(() => fetchFn(url, init), opts); + return async (url, init = {}) => { + const { retrySafe, ...request } = init; + const method = String(request.method || 'GET').toUpperCase(); + const safe = retrySafe === true || opts.retrySafe === true || ['GET', 'HEAD', 'OPTIONS'].includes(method); + const timeoutMs = opts.timeoutMs ?? envInt('ATS_HTTP_TIMEOUT_MS', 30_000); + if (!Number.isFinite(timeoutMs) || timeoutMs <= 0) throw new Error('HTTP timeoutMs must be a positive number'); + const controller = new AbortController(); + const signals = [request.signal, opts.signal, controller.signal].filter(Boolean); + const signal = signals.length === 1 ? signals[0] : globalThis.AbortSignal.any(signals); + const timer = setTimeout(() => controller.abort(Object.assign( + new Error(`HTTP request timed out after ${timeoutMs}ms`), { name: 'TimeoutError', code: 'ATS_TIMEOUT', exitCode: 6 } + )), timeoutMs); + let response; + try { + response = await withRetry(() => fetchFn(url, { ...request, signal }), { + ...opts, signal, + isRetryableError: (error) => safe && (opts.isRetryableError ? opts.isRetryableError(error) : isTransientError(error)), + isRetryableResponse: (res, body) => safe + ? (opts.isRetryableResponse ? opts.isRetryableResponse(res, body) : isTransientResponse(res, body)) + : res?.status === 429 || (res?.status === 403 && upstreamWaitMs(res) !== null) || (res?.status === 500 && EXPLICIT_RATE_LIMIT_BODY.test(body)), + }); + } catch (error) { + clearTimeout(timer); + throw error; + } + // The same deadline covers the response body, not just response headers. + // Unconsumed responses do not keep a CLI process alive solely for this timer. + timer.unref?.(); + for (const method of ['text', 'json', 'arrayBuffer', 'blob', 'formData']) { + if (typeof response?.[method] !== 'function') continue; + const read = response[method].bind(response); + response[method] = async (...args) => { + timer.ref?.(); + try { return await abortable(read(...args), signal); } + finally { clearTimeout(timer); } + }; + } + return response; + }; } diff --git a/packages/core/test/kg-integrity.test.js b/packages/core/test/kg-integrity.test.js new file mode 100644 index 0000000..93701df --- /dev/null +++ b/packages/core/test/kg-integrity.test.js @@ -0,0 +1,67 @@ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { spawn } from 'node:child_process'; +import { proposeFact, ratifyFactItem, loadFacts, proposeConfirm, staleFacts, exportFactsCypher, exportFactsGraphiti } from '../kg.js'; +import { decideReviewItem } from '../review-queue.js'; + +function store(t) { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'ats-kg-integrity-')); + t.after(() => fs.rmSync(dir, { recursive: true, force: true })); + return { factsPath: path.join(dir, 'facts.jsonl'), queuePath: path.join(dir, 'review.json') }; +} +const input = { subject: 'Acme', predicate: 'uses', object: 'Invoices', by: 'agent' }; +const approve = (item, paths) => decideReviewItem(item.id, 'approve', { ...paths, by: 'reviewer' }); + +test('ratification refuses tampered approval and repeated concurrent ratification appends once', async (t) => { + const paths = store(t); + const item = approve(proposeFact(input, paths), paths); + assert.throws(() => ratifyFactItem({ ...item, payload: { ...item.payload, object: 'changed' } }, paths), /matching payload approval/); + const source = `import { ratifyFactItem } from ${JSON.stringify(new URL('../kg.js', import.meta.url).href)};ratifyFactItem(${JSON.stringify(item)},${JSON.stringify(paths)});`; + await Promise.all(Array.from({ length: 6 }, () => new Promise((resolve, reject) => { + const child = spawn(process.execPath, ['--input-type=module', '-e', source]); + let err = ''; child.stderr.on('data', (chunk) => { err += chunk; }); + child.on('error', reject); child.on('exit', (code) => code === 0 ? resolve() : reject(new Error(err))); + }))); + assert.equal(loadFacts(paths).facts.length, 1); + assert.equal(loadFacts(paths).events, 1); +}); + +test('a conflicting fact staged before another ratification cannot become silently current', (t) => { + const paths = store(t); + const a = approve(proposeFact(input, paths), paths); + const b = approve(proposeFact({ ...input, object: 'Other' }, paths), paths); + ratifyFactItem(a, paths); + assert.throws(() => ratifyFactItem(b, paths), /active fact/); + assert.equal(loadFacts(paths).events, 1); +}); + +test('valid and learned dates survive ratification and export; future dates are refused', (t) => { + const paths = store(t); + assert.throws(() => proposeFact({ ...input, validAt: '2999-01-01T00:00:00Z' }, paths), /future/); + assert.throws(() => proposeFact({ ...input, learnedAt: '2026-02-30T00:00:00Z' }, paths), /invalid/); + const item = approve(proposeFact({ ...input, validAt: '2026-01-01T02:00:00+02:00', learnedAt: '2026-02-01T00:00:00Z' }, paths), paths); + const fact = ratifyFactItem(item, { ...paths, now: '2026-03-01T00:00:00Z' }).fact; + assert.equal(fact.tValid, '2026-01-01T00:00:00.000Z'); + assert.equal(fact.tLearned, '2026-02-01T00:00:00.000Z'); + assert.equal(fact.provenance.ratifiedAt, '2026-03-01T00:00:00.000Z'); + assert.match(exportFactsCypher(paths), /tLearned/); + assert.equal(JSON.parse(exportFactsGraphiti(paths)).learned_at, fact.tLearned); +}); + +test('stale facts need reviewed confirmation; original validity and provenance remain intact', (t) => { + const paths = store(t); + const fact = ratifyFactItem(approve(proposeFact(input, paths), paths), { ...paths, now: '2026-01-01T00:00:00Z' }).fact; + const now = Date.parse('2026-04-01T00:00:00Z'); + assert.equal(staleFacts({ ...paths, days: 60, now }).count, 1); + const confirm = proposeConfirm({ factId: fact.id, source: 'read-only source check', by: 'agent' }, paths); + assert.throws(() => ratifyFactItem(confirm, paths), /not approved/); + ratifyFactItem(approve(confirm, paths), { ...paths, now: '2026-03-31T00:00:00Z' }); + assert.equal(staleFacts({ ...paths, days: 60, now }).count, 0); + const fresh = loadFacts(paths).facts[0]; + assert.equal(fresh.tValid, fact.tValid); + assert.deepEqual(fresh.provenance, fact.provenance); + assert.equal(fresh.confirmation.source, 'read-only source check'); +}); diff --git a/packages/core/test/retry.test.js b/packages/core/test/retry.test.js index 525e0b8..893e34e 100644 --- a/packages/core/test/retry.test.js +++ b/packages/core/test/retry.test.js @@ -133,3 +133,55 @@ test('retryingFetch wraps a fetch function transparently and reports retries', a assert.equal(events.length, 1); assert.match(events[0].reason, /Demo: HTTP 503/); }); + + +test('fetch deadline aborts a stalled request and body; backoff obeys cancellation', async () => { + let signal; + const stalled = retryingFetch(async (_url, init) => { signal = init.signal; return new Promise(() => {}); }, { timeoutMs: 15 }); + await assert.rejects(stalled('https://example.test'), (error) => error.code === 'ATS_TIMEOUT'); + assert.equal(signal.aborted, true); + const body = retryingFetch(async () => ({ status: 200, text: () => new Promise(() => {}) }), { timeoutMs: 15 }); + const response = await body('https://example.test'); + await assert.rejects(response.text(), (error) => error.code === 'ATS_TIMEOUT'); + + const controller = new AbortController(); + let calls = 0; + const retry = retryingFetch(async () => { calls++; return resp(429, { headers: { 'retry-after': '50' } }); }, { + signal: controller.signal, onRetry: () => controller.abort(new Error('cancelled by caller')), + }); + await assert.rejects(retry('https://example.test'), /cancelled by caller/); + assert.equal(calls, 1); +}); + +test('unsafe mutations are not replayed on gateway or network ambiguity; read-only POST opts in', async () => { + let calls = 0; + const gateway = retryingFetch(async () => { calls++; return resp(503); }, { sleep: noSleep }); + assert.equal((await gateway('https://example.test', { method: 'POST' })).status, 503); + assert.equal(calls, 1); + calls = 0; + const ambiguous = retryingFetch(async () => { calls++; return resp(500, { body: 'temporarily unavailable; try again later' }); }, { sleep: noSleep }); + assert.equal((await ambiguous('https://example.test', { method: 'POST' })).status, 500); + assert.equal(calls, 1, 'generic retry language does not prove a rejected write'); + calls = 0; + const rateLimited = retryingFetch(async () => ++calls === 1 ? resp(500, { body: 'exceed_query_limit' }) : resp(200), { sleep: noSleep }); + assert.equal((await rateLimited('https://example.test', { method: 'POST' })).status, 200); + assert.equal(calls, 2); + calls = 0; + const broken = retryingFetch(async () => { calls++; throw new TypeError('fetch failed'); }, { sleep: noSleep }); + await assert.rejects(broken('https://example.test', { method: 'PATCH' }), /fetch failed/); + assert.equal(calls, 1); + calls = 0; + const read = retryingFetch(async (_url, init) => { + assert.equal(init.retrySafe, undefined, 'ATS option is not sent to fetch'); + return ++calls === 1 ? resp(503) : resp(200); + }, { sleep: noSleep }); + assert.equal((await read('https://example.test/query', { method: 'POST', retrySafe: true })).status, 200); + assert.equal(calls, 2); +}); + +test('retry does not send early when the server asks for a long reset', async () => { + let calls = 0; + const fetch = retryingFetch(async () => { calls++; return resp(429, { headers: { 'retry-after': '120' } }); }, { sleep: noSleep }); + assert.equal((await fetch('https://example.test')).status, 429); + assert.equal(calls, 1); +}); diff --git a/packages/mcp/package.json b/packages/mcp/package.json index f1e1725..1704706 100644 --- a/packages/mcp/package.json +++ b/packages/mcp/package.json @@ -1,6 +1,6 @@ { "name": "@reneza/ats-mcp", - "version": "0.15.0", + "version": "0.16.0", "mcpName": "io.github.renezander030/agentic-task-system", "description": "Model Context Protocol server for Agentic Task System — exposes the task app you already use to any MCP client, backed by hybrid + RRF retrieval. Storage-agnostic over the ATS adapter contract.", "type": "module", @@ -20,11 +20,11 @@ }, "dependencies": { "@modelcontextprotocol/sdk": "^1.30.1", - "@reneza/ats-core": "^0.15.0", + "@reneza/ats-core": "^0.16.0", "zod": "^4.6.5" }, "peerDependencies": { - "@reneza/ats-adapter-ticktick": "^0.15.0" + "@reneza/ats-adapter-ticktick": "^0.16.0" }, "peerDependenciesMeta": { "@reneza/ats-adapter-ticktick": { diff --git a/server.json b/server.json index ad7c0af..2ad5b1e 100644 --- a/server.json +++ b/server.json @@ -1,7 +1,7 @@ { "$schema": "https://static.modelcontextprotocol.io/schemas/2025-12-11/server.schema.json", "name": "io.github.renezander030/agentic-task-system", - "description": "MCP server giving AI agents persistent task memory across TickTick, Notion, GitHub, Linear & Beads", + "description": "Persistent task context across TickTick, Notion, GitHub, Obsidian, Taskmaster and Beads", "title": "Agentic Task System (ATS)", "websiteUrl": "https://github.com/renezander030/agentic-task-system", "repository": { @@ -9,13 +9,13 @@ "source": "github", "subfolder": "packages/mcp" }, - "version": "0.15.0", + "version": "0.16.0", "packages": [ { "registryType": "npm", "registryBaseUrl": "https://registry.npmjs.org", "identifier": "@reneza/ats-mcp", - "version": "0.15.0", + "version": "0.16.0", "transport": { "type": "stdio" },