Compare commits

...

3 Commits

Author SHA1 Message Date
david 52c4d1f1e6 fix: derived-state determinism, error visibility, metrics naming, HTTP hardening
- rebuild: fold observations per fingerprint in a total order
  (hcl DESC, node_id DESC, id DESC) so two nodes with identical logs
  rebuild identical entries (was: arbitrary bare-column row, merge-order
  dependent)
- threads: CreateThreadCluster uses INSERT OR IGNORE + existing-id
  fallback — concurrent auto-threaders converge instead of hitting UNIQUE
- errors surfaced instead of swallowed: scanEntries returns rows.Err(),
  Stats() fails fast on query errors, AutoLinkThreadObservations /
  linkTemporalNeighbors / golden-thread linking propagate failures,
  AddThreadNote + LinkObservationToThread write under one tx,
  watch records session upsert failures
- gossip client: push checks HTTP status and reports errors (a broken
  push direction no longer looks like a silent success); gossip diff
  gets a 10s timeout so a dead peer cannot hang the CLI
- ingest: failed source sweeps (obsidian/browser/gitea) are logged,
  and -d's help text now states its file-only scope
- watch --quiet: fatal errors go to stderr instead of io.Discard
- main: cobra SilenceErrors/SilenceUsage (errors print once, usage is
  not dumped on runtime failures); knox mcp exits 0 on SIGINT/SIGTERM
- metrics: drop _total suffix from gauges (knox_observations,
  knox_entries, knox_projects, knox_sessions, knox_peers, knox_threads);
  _total stays on counters per Prometheus convention
- http: ReadHeaderTimeout + IdleTimeout on gossip, metrics, and web servers

tests: concurrent cluster-create idempotency, HCL-order rebuild fold
(both merge orders), push HTTP-error surfacing; full suite + -race pass,
gofmt clean
2026-09-17 02:28:08 -07:00
david bb852faa27 fix: harden gossip, HLC restarts, watcher races, MCP args, pagination (#3)
Implements the top findings from the codebase review, verified with tests and live CLI/MCP checks.

**Gossip integrity**
- Push validation: 4 MiB body cap, 1000-row batch cap; rows claiming the local node id (vector-poisoning), empty node ids, and negative HCLs rejected (internal/watch/gossip.go, internal/db/gossip.go)
- Reconcile-on-pull: Run returns the pulled count, syncGossip rebuilds derived state when > 0 — entry-count comparison could never fire, so synced observations never materialized into searchable entries

**Data-layer safety**
- HLC resumed from MAX(hcl) at Open (hlc.SeekTo): a restart with a regressed wall clock cannot reissue values the (node_id, hcl) locator and pull cursors depend on
- Writer serialization: _txlock=immediate DSN + SetMaxOpenConns(1) + per-KnoxDB mutex around RecordObservation's check-then-insert dedup (closes duplicate-row race)

**Watch daemon**
- Ticker guard flags now atomic.Bool (was a cross-goroutine data race)
- Trailing-edge per-path debounce (timer-based, pruned on fire/delete)
- Recursive watches (startup tree walk + watcher.Add on dir Create); Rename re-ingests, Remove cancels pending ingests

**MCP + CLI**
- Strict arg validation, no silent clamping: thread_id 0 errors instead of renaming thread #1; empty knox_thread_link {} errors instead of false success; thread existence checked before writes; golden-thread tool nil-safe
- --page 0 errors instead of panicking; query/recent pagination actually pages (page x limit)

**Tests** (new internal/hlc and internal/db packages): SeekTo monotonicity, concurrent dedup race, push validation, reopen HCL monotonicity, batch caps, self-spoof rejection, idempotency on observation counts.

Verified: go build, go vet, full suite with -race, live MCP stdio transcripts against a scratch DB.
Reviewed-on: #3
Co-authored-by: David Gwilliam <dhgwilliam@gmail.com>
Co-committed-by: David Gwilliam <dhgwilliam@gmail.com>
2026-09-17 09:06:08 +00:00
david d6d2a24ddc style: gofmt all packages (#4)
Applies gofmt to the 18 files that were already unformatted at HEAD (pre-existing debt — 122 insertions / 122 deletions, whitespace plus import-block reorderings only; `git diff -w` confirms no semantic changes).

Kept on its own branch so the functional change set (see the review-hardening PR) stays reviewable without formatting noise.

Verified: go build, go vet, go test ./... pass on this branch; a merge simulation with the functional branch produces a clean 3-way merge with all tests green.
Reviewed-on: #4
Co-authored-by: David Gwilliam <dhgwilliam@gmail.com>
Co-committed-by: David Gwilliam <dhgwilliam@gmail.com>
2026-09-17 09:05:57 +00:00
24 changed files with 1152 additions and 329 deletions
+5 -5
View File
@@ -5,9 +5,9 @@ import (
"log" "log"
"strings" "strings"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/ingest" "github.com/david/knox/internal/ingest"
"github.com/spf13/cobra"
) )
func NewBrowserCmd(kdb *db.KnoxDB) *cobra.Command { func NewBrowserCmd(kdb *db.KnoxDB) *cobra.Command {
@@ -65,10 +65,10 @@ Run periodically to capture browsing exhaust.`,
for _, r := range results { for _, r := range results {
domains[extractDomain(r.Summary)]++ domains[extractDomain(r.Summary)]++
} }
fmt.Println("\nTop domains:") fmt.Println("\nTop domains:")
for _, d := range topDomains(domains, 10) { for _, d := range topDomains(domains, 10) {
fmt.Printf(" %-40s %d\n", d.Name, d.Count) fmt.Printf(" %-40s %d\n", d.Name, d.Count)
} }
return nil return nil
}, },
+1 -1
View File
@@ -4,9 +4,9 @@ import (
"fmt" "fmt"
"log" "log"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/ingest" "github.com/david/knox/internal/ingest"
"github.com/spf13/cobra"
) )
func NewGiteaCmd(kdb *db.KnoxDB) *cobra.Command { func NewGiteaCmd(kdb *db.KnoxDB) *cobra.Command {
+3 -2
View File
@@ -2,10 +2,11 @@ package cmd
import ( import (
"fmt" "fmt"
"time"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/watch" "github.com/david/knox/internal/watch"
"github.com/spf13/cobra"
) )
// NewGossipCmd exposes peer status/diff for the gossip protocol. // NewGossipCmd exposes peer status/diff for the gossip protocol.
@@ -77,7 +78,7 @@ func NewGossipCmd(kdb *db.KnoxDB) *cobra.Command {
// gossipDiff compares this node's observation log + auto-threads with a peer's // gossipDiff compares this node's observation log + auto-threads with a peer's
// via /v1/diff. Never writes; it is the preview for what a sync would adopt. // via /v1/diff. Never writes; it is the preview for what a sync would adopt.
func gossipDiff(kdb *db.KnoxDB, peerAddr string) error { func gossipDiff(kdb *db.KnoxDB, peerAddr string) error {
c := &watch.Client{Addr: peerAddr} c := &watch.Client{Addr: peerAddr, Timeout: 10 * time.Second}
remote, err := c.Diff() remote, err := c.Diff()
if err != nil { if err != nil {
return err return err
+79 -73
View File
@@ -5,10 +5,10 @@ import (
"log" "log"
"path/filepath" "path/filepath"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/ingest" "github.com/david/knox/internal/ingest"
"github.com/david/knox/internal/watch" "github.com/david/knox/internal/watch"
"github.com/spf13/cobra"
) )
func NewIngestCmd(kdb *db.KnoxDB) *cobra.Command { func NewIngestCmd(kdb *db.KnoxDB) *cobra.Command {
@@ -27,92 +27,98 @@ func NewIngestCmd(kdb *db.KnoxDB) *cobra.Command {
ingest.NewSkillsIngester(), ingest.NewSkillsIngester(),
} }
var totalNew, totalUpdated int var totalNew, totalUpdated int
record := func(result *ingest.IngestResult) { record := func(result *ingest.IngestResult) {
_, isNew, err := kdb.RecordObservation(db.ObservationRecord{ _, isNew, err := kdb.RecordObservation(db.ObservationRecord{
Fingerprint: result.Fingerprint, Fingerprint: result.Fingerprint,
SourceID: result.SourceID, SourceID: result.SourceID,
SourcePath: result.SourcePath, SourcePath: result.SourcePath,
Project: result.Project, Project: result.Project,
ContentType: result.ContentType, ContentType: result.ContentType,
Title: result.Title, Title: result.Title,
Summary: result.Summary, Summary: result.Summary,
CreatedAt: result.CreatedAt, CreatedAt: result.CreatedAt,
LineStart: result.LineStart, LineStart: result.LineStart,
LineEnd: result.LineEnd, LineEnd: result.LineEnd,
Confidence: result.Confidence, Confidence: result.Confidence,
IngesterVersion: result.IngesterVersion, IngesterVersion: result.IngesterVersion,
Trigger: "ingest", Trigger: "ingest",
Provenance: ingest.ProvenanceJSON(result.Provenance), Provenance: ingest.ProvenanceJSON(result.Provenance),
}) })
if err != nil { if err != nil {
log.Printf("[knox] db error: %v", err) log.Printf("[knox] db error: %v", err)
return return
}
if isNew {
totalNew++
} else {
totalUpdated++
}
if result.SourceID == "opencode-session" {
sessionID, _ := result.Provenance["session_id"].(string)
if sessionID != "" {
kdb.UpsertSession(sessionID, result.Project, result.Title, "active")
} }
} if isNew {
} totalNew++
} else {
for _, ing := range ingesters { totalUpdated++
sid := ing.SourceID()
log.Printf("[knox] ingesting %s...", sid)
for _, dir := range dirs {
patterns := []string{
dir + "/*",
dir + "/*/SKILL.md",
} }
for _, pattern := range patterns {
entries, _ := filepath.Glob(pattern)
for _, path := range entries {
if !watch.MatchesIngester(path, sid) {
continue
}
result, err := ing.Ingest(path) if result.SourceID == "opencode-session" {
if err != nil || result == nil { sessionID, _ := result.Provenance["session_id"].(string)
continue if sessionID != "" {
} kdb.UpsertSession(sessionID, result.Project, result.Title, "active")
record(result)
} }
} }
} }
}
// Non-file sources: obsidian vault, browser history, gitea for _, ing := range ingesters {
if vault, err := ingest.DetectObsidianVault(); err == nil { sid := ing.SourceID()
log.Printf("[knox] ingesting obsidian...") log.Printf("[knox] ingesting %s...", sid)
if results, err := ingest.NewObsidianIngester(vault).IngestAll(); err == nil {
for _, dir := range dirs {
patterns := []string{
dir + "/*",
dir + "/*/SKILL.md",
}
for _, pattern := range patterns {
entries, _ := filepath.Glob(pattern)
for _, path := range entries {
if !watch.MatchesIngester(path, sid) {
continue
}
result, err := ing.Ingest(path)
if err != nil || result == nil {
continue
}
record(result)
}
}
}
}
// Non-file sources: obsidian vault, browser history, gitea
if vault, err := ingest.DetectObsidianVault(); err == nil {
log.Printf("[knox] ingesting obsidian...")
if results, err := ingest.NewObsidianIngester(vault).IngestAll(); err == nil {
for _, r := range results {
record(r)
}
} else {
log.Printf("[knox] obsidian ingest failed: %v", err)
}
}
log.Printf("[knox] ingesting browser-history...")
if results, err := ingest.NewBrowserHistoryIngester().IngestAll(); err == nil {
for _, r := range results { for _, r := range results {
record(r) record(r)
} }
} else {
log.Printf("[knox] browser-history ingest failed: %v", err)
} }
} log.Printf("[knox] ingesting gitea...")
log.Printf("[knox] ingesting browser-history...") if results, err := ingest.NewGiteaIngester().IngestAll(); err == nil {
if results, err := ingest.NewBrowserHistoryIngester().IngestAll(); err == nil { for _, r := range results {
for _, r := range results { record(r)
record(r) }
} else {
log.Printf("[knox] gitea ingest failed: %v", err)
} }
}
log.Printf("[knox] ingesting gitea...")
if results, err := ingest.NewGiteaIngester().IngestAll(); err == nil {
for _, r := range results {
record(r)
}
}
fmt.Printf("\nIngest complete: %d new entries, %d updated\n", totalNew, totalUpdated) fmt.Printf("\nIngest complete: %d new entries, %d updated\n", totalNew, totalUpdated)
}, },
} }
cmd.Flags().StringSliceVarP(&dirs, "dir", "d", nil, "Directories to scan (default: opencode storage/log + skills)") cmd.Flags().StringSliceVarP(&dirs, "dir", "d", nil, "Directories to scan (default: opencode storage/log + skills)")
+180 -40
View File
@@ -4,6 +4,7 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"fmt" "fmt"
"math"
"strings" "strings"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
@@ -50,8 +51,14 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(searchTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(searchTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
query, _ := req.Params.Arguments["query"].(string) query, err := requiredStringArg(req, "query")
limit := clampInt(req, "limit", 10, 1, 50) if err != nil {
return errorResult(err.Error()), nil
}
limit, err := optionalIntArg(req, "limit", 10, 1, 50)
if err != nil {
return errorResult(err.Error()), nil
}
detail := parseDetail(getString(req, "detail", "normal")) detail := parseDetail(getString(req, "detail", "normal"))
scope := getString(req, "scope", "all") scope := getString(req, "scope", "all")
@@ -76,7 +83,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(recentTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(recentTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
limit := clampInt(req, "limit", 10, 1, 50) limit, err := optionalIntArg(req, "limit", 10, 1, 50)
if err != nil {
return errorResult(err.Error()), nil
}
detail := parseDetail(getString(req, "detail", "normal")) detail := parseDetail(getString(req, "detail", "normal"))
scope := getString(req, "scope", "all") scope := getString(req, "scope", "all")
@@ -99,7 +109,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(getTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(getTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
fp := getString(req, "fingerprint", "") fp, err := requiredStringArg(req, "fingerprint")
if err != nil {
return errorResult(err.Error()), nil
}
entry, err := kdb.FindEntry(fp) entry, err := kdb.FindEntry(fp)
if err != nil { if err != nil {
@@ -165,7 +178,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(countTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(countTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
query, _ := req.Params.Arguments["query"].(string) query, err := requiredStringArg(req, "query")
if err != nil {
return errorResult(err.Error()), nil
}
results, err := kdb.Search(query, 500) results, err := kdb.Search(query, 500)
if err != nil { if err != nil {
return errorResult(err.Error()), nil return errorResult(err.Error()), nil
@@ -235,7 +251,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(threadCreateTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(threadCreateTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
title := getString(req, "title", "") title, err := requiredStringArg(req, "title")
if err != nil {
return errorResult(err.Error()), nil
}
motivation := getString(req, "motivation", "") motivation := getString(req, "motivation", "")
priority := getString(req, "priority", "medium") priority := getString(req, "priority", "medium")
provenance := getString(req, "provenance", "{}") provenance := getString(req, "provenance", "{}")
@@ -257,7 +276,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(threadUpdateTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(threadUpdateTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
id := int64(clampInt(req, "thread_id", 0, 1, 999999)) id, err := requiredIntArg(req, "thread_id", 1, 999999)
if err != nil {
return errorResult(err.Error()), nil
}
title := getString(req, "title", "") title := getString(req, "title", "")
motivation := getString(req, "motivation", "") motivation := getString(req, "motivation", "")
priority := getString(req, "priority", "") priority := getString(req, "priority", "")
@@ -281,7 +303,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(threadDraftTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(threadDraftTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
limit := clampInt(req, "limit", 20, 1, 200) limit, err := optionalIntArg(req, "limit", 20, 1, 200)
if err != nil {
return errorResult(err.Error()), nil
}
threads, err := kdb.ListThreads("") threads, err := kdb.ListThreads("")
if err != nil { if err != nil {
return errorResult(err.Error()), nil return errorResult(err.Error()), nil
@@ -329,9 +354,18 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(threadLinkTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(threadLinkTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
threadID := int64(clampInt(req, "thread_id", 0, 1, 999999)) threadID, err := requiredIntArg(req, "thread_id", 1, 999999)
obsID := int64(clampInt(req, "observation_id", 0, 1, 999999)) if err != nil {
return errorResult(err.Error()), nil
}
obsID, err := requiredIntArg(req, "observation_id", 1, 999999)
if err != nil {
return errorResult(err.Error()), nil
}
relevance := getString(req, "relevance", "") relevance := getString(req, "relevance", "")
if err := threadExists(kdb, threadID); err != nil {
return errorResult(err.Error()), nil
}
if err := kdb.LinkObservationToThread(threadID, obsID, relevance); err != nil { if err := kdb.LinkObservationToThread(threadID, obsID, relevance); err != nil {
return errorResult(err.Error()), nil return errorResult(err.Error()), nil
} }
@@ -345,7 +379,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(threadSearchTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(threadSearchTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
query := getString(req, "query", "") query, err := requiredStringArg(req, "query")
if err != nil {
return errorResult(err.Error()), nil
}
threads, err := kdb.SearchThreadsByMotivation(query) threads, err := kdb.SearchThreadsByMotivation(query)
if err != nil { if err != nil {
return errorResult(err.Error()), nil return errorResult(err.Error()), nil
@@ -372,9 +409,18 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(threadLinkEntryTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(threadLinkEntryTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
threadID := int64(clampInt(req, "thread_id", 0, 1, 999999)) threadID, err := requiredIntArg(req, "thread_id", 1, 999999)
fp := getString(req, "fingerprint", "") if err != nil {
return errorResult(err.Error()), nil
}
fp, err := requiredStringArg(req, "fingerprint")
if err != nil {
return errorResult(err.Error()), nil
}
relation := getString(req, "relation", "produced") relation := getString(req, "relation", "produced")
if err := threadExists(kdb, threadID); err != nil {
return errorResult(err.Error()), nil
}
if err := kdb.LinkEntryToThread(threadID, fp, relation); err != nil { if err := kdb.LinkEntryToThread(threadID, fp, relation); err != nil {
return errorResult(err.Error()), nil return errorResult(err.Error()), nil
} }
@@ -390,9 +436,21 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(threadRelateTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(threadRelateTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
parentID := int64(clampInt(req, "parent_id", 0, 1, 999999)) parentID, err := requiredIntArg(req, "parent_id", 1, 999999)
childID := int64(clampInt(req, "child_id", 0, 1, 999999)) if err != nil {
return errorResult(err.Error()), nil
}
childID, err := requiredIntArg(req, "child_id", 1, 999999)
if err != nil {
return errorResult(err.Error()), nil
}
relation := getString(req, "relation", "spawned") relation := getString(req, "relation", "spawned")
if err := threadExists(kdb, parentID); err != nil {
return errorResult(err.Error()), nil
}
if err := threadExists(kdb, childID); err != nil {
return errorResult(err.Error()), nil
}
if err := kdb.LinkEntryToThread(childID, db.ThreadFP(parentID), relation); err != nil { if err := kdb.LinkEntryToThread(childID, db.ThreadFP(parentID), relation); err != nil {
return errorResult(err.Error()), nil return errorResult(err.Error()), nil
} }
@@ -406,29 +464,67 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(goldenTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(goldenTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
threadID := int64(clampInt(req, "thread_id", 0, 0, 999999)) // Absent thread_id: query the current golden thread.
if _, ok := req.Params.Arguments["thread_id"]; !ok {
if threadID == 0 { current, err := kdb.GoldenThreadID()
// Query or clear if err != nil {
current, _ := kdb.GoldenThreadID() return errorResult(err.Error()), nil
}
if current == 0 { if current == 0 {
return mcp.NewToolResultText("No golden thread set."), nil return mcp.NewToolResultText("No golden thread set."), nil
} }
// If thread_id was explicitly 0 and a golden exists, clear it t, err := kdb.GetThread(current)
if _, ok := req.Params.Arguments["thread_id"]; ok { if err != nil {
kdb.SetGoldenThread(0) return errorResult(err.Error()), nil
t, _ := kdb.GetThread(current) }
return mcp.NewToolResultText(fmt.Sprintf("Golden thread cleared (was #%d: %s)", current, t.Title)), nil if t == nil {
return errorResult(fmt.Sprintf("Golden thread #%d no longer exists", current)), nil
} }
t, _ := kdb.GetThread(current)
return mcp.NewToolResultText(fmt.Sprintf("Golden thread: #%d %s [%s]\n %s", t.ID, t.Title, t.Status, t.Motivation)), nil return mcp.NewToolResultText(fmt.Sprintf("Golden thread: #%d %s [%s]\n %s", t.ID, t.Title, t.Status, t.Motivation)), nil
} }
threadID, err := requiredIntArg(req, "thread_id", 0, 999999)
if err != nil {
return errorResult(err.Error()), nil
}
if threadID == 0 {
// Explicit 0 clears the golden thread.
current, err := kdb.GoldenThreadID()
if err != nil {
return errorResult(err.Error()), nil
}
if current == 0 {
return mcp.NewToolResultText("No golden thread set."), nil
}
if err := kdb.SetGoldenThread(0); err != nil {
return errorResult(err.Error()), nil
}
t, err := kdb.GetThread(current)
if err != nil {
return errorResult(err.Error()), nil
}
name := "?"
if t != nil {
name = t.Title
}
return mcp.NewToolResultText(fmt.Sprintf("Golden thread cleared (was #%d: %s)", current, name)), nil
}
if err := threadExists(kdb, threadID); err != nil {
return errorResult(err.Error()), nil
}
if err := kdb.SetGoldenThread(threadID); err != nil { if err := kdb.SetGoldenThread(threadID); err != nil {
return errorResult(err.Error()), nil return errorResult(err.Error()), nil
} }
t, _ := kdb.GetThread(threadID) t, err := kdb.GetThread(threadID)
return mcp.NewToolResultText(fmt.Sprintf("Golden thread set to #%d: %s", threadID, t.Title)), nil if err != nil {
return errorResult(err.Error()), nil
}
name := "?"
if t != nil {
name = t.Title
}
return mcp.NewToolResultText(fmt.Sprintf("Golden thread set to #%d: %s", threadID, name)), nil
}) })
// ─── knox_topics ─────────────────────────────────────────── // ─── knox_topics ───────────────────────────────────────────
@@ -438,7 +534,10 @@ func NewMCPServer(kdb *db.KnoxDB) *server.MCPServer {
) )
s.AddTool(topicsTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { s.AddTool(topicsTool, func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
limit := clampInt(req, "limit", 10, 1, 30) limit, err := optionalIntArg(req, "limit", 10, 1, 30)
if err != nil {
return errorResult(err.Error()), nil
}
entries, _ := kdb.RecentEntries(2000) entries, _ := kdb.RecentEntries(2000)
tfidf := index.BuildTFIDF(entries) tfidf := index.BuildTFIDF(entries)
clusters := tfidf.Cluster(2, limit) clusters := tfidf.Cluster(2, limit)
@@ -618,18 +717,59 @@ func getString(req mcp.CallToolRequest, key, def string) string {
return def return def
} }
func clampInt(req mcp.CallToolRequest, key string, def, min, max int) int { // requiredStringArg returns a non-empty string argument or a descriptive error.
if v, ok := req.Params.Arguments[key].(float64); ok { // Empty/missing/wrong-type values are rejected rather than silently defaulted:
n := int(v) // LLM clients routinely omit or zero value params, and a silent default writes
if n < min { // to the wrong thread or reports false success.
return min func requiredStringArg(req mcp.CallToolRequest, key string) (string, error) {
} v, _ := req.Params.Arguments[key].(string)
if n > max { v = strings.TrimSpace(v)
return max if v == "" {
} return "", fmt.Errorf("%s is required and must be a non-empty string", key)
return n
} }
return def return v, nil
}
// requiredIntArg returns an integer argument validated against [min, max].
func requiredIntArg(req mcp.CallToolRequest, key string, min, max int64) (int64, error) {
v, ok := req.Params.Arguments[key].(float64)
if !ok || v != math.Trunc(v) {
return 0, fmt.Errorf("%s is required and must be an integer", key)
}
n := int64(v)
if n < min || n > max {
return 0, fmt.Errorf("%s must be in [%d..%d], got %d", key, min, max, n)
}
return n, nil
}
// optionalIntArg validates a present numeric argument, defaulting when absent.
func optionalIntArg(req mcp.CallToolRequest, key string, def, min, max int) (int, error) {
v, ok := req.Params.Arguments[key].(float64)
if !ok {
return def, nil
}
if v != math.Trunc(v) {
return 0, fmt.Errorf("%s must be an integer", key)
}
n := int(v)
if n < min || n > max {
return 0, fmt.Errorf("%s must be in [%d..%d], got %d", key, min, max, n)
}
return n, nil
}
// threadExists is a guard for write tools: link/relate/set-golden must not
// silently accept ids with no matching thread.
func threadExists(kdb *db.KnoxDB, id int64) error {
t, err := kdb.GetThread(id)
if err != nil {
return err
}
if t == nil {
return fmt.Errorf("thread #%d not found", id)
}
return nil
} }
func shortFP(fp string) string { func shortFP(fp string) string {
+1 -1
View File
@@ -4,9 +4,9 @@ import (
"fmt" "fmt"
"log" "log"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/ingest" "github.com/david/knox/internal/ingest"
"github.com/spf13/cobra"
) )
func NewObsidianCmd(kdb *db.KnoxDB) *cobra.Command { func NewObsidianCmd(kdb *db.KnoxDB) *cobra.Command {
+1 -1
View File
@@ -3,9 +3,9 @@ package cmd
import ( import (
"fmt" "fmt"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/watch" "github.com/david/knox/internal/watch"
"github.com/spf13/cobra"
) )
// NewReconcileCmd rebuilds all derived state (entries cache, threads) from the // NewReconcileCmd rebuilds all derived state (entries cache, threads) from the
+20 -7
View File
@@ -6,14 +6,14 @@ import (
"os" "os"
"strings" "strings"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/index" "github.com/david/knox/internal/index"
"github.com/spf13/cobra"
) )
func paginate(entries []db.Entry, page, limit int) []db.Entry { func paginate(entries []db.Entry, page, limit int) []db.Entry {
start := (page - 1) * limit start := (page - 1) * limit
if start >= len(entries) { if start < 0 || start >= len(entries) {
return nil return nil
} }
end := start + limit end := start + limit
@@ -74,16 +74,22 @@ func NewQueryCmd(kdb *db.KnoxDB) *cobra.Command {
if query == "" { if query == "" {
return fmt.Errorf("search query required (positional arg or stdin pipe)") return fmt.Errorf("search query required (positional arg or stdin pipe)")
} }
results, err := kdb.Search(query, limit) if page < 1 || limit < 1 {
return fmt.Errorf("page and limit must be >= 1 (got page=%d limit=%d)", page, limit)
}
// Fetch enough rows to cover the requested page: Search caps at its
// limit argument, so slicing a limit-sized result set could never
// reach page 2+.
all, err := kdb.Search(query, page*limit)
if err != nil { if err != nil {
return err return err
} }
results = paginate(results, page, limit) results := paginate(all, page, limit)
if len(results) == 0 { if len(results) == 0 {
fmt.Println("No results found.") fmt.Println("No results found.")
return nil return nil
} }
fmt.Printf("Found %d results for %q (page %d, %d per page):\n\n", len(results), query, page, limit) fmt.Printf("Found %d results for %q — showing %d (page %d, %d per page):\n\n", len(all), query, len(results), page, limit)
for _, r := range results { for _, r := range results {
fmt.Printf(" %-8s %-30s [%s] %s\n", shortFP(r.Fingerprint), truncateStr(r.Title, 30), r.Project, r.SourceID) fmt.Printf(" %-8s %-30s [%s] %s\n", shortFP(r.Fingerprint), truncateStr(r.Title, 30), r.Project, r.SourceID)
if r.Summary != "" { if r.Summary != "" {
@@ -105,13 +111,20 @@ func NewRecentCmd(kdb *db.KnoxDB) *cobra.Command {
Use: "recent", Use: "recent",
Short: "Show recent knowledge entries", Short: "Show recent knowledge entries",
RunE: func(c *cobra.Command, args []string) error { RunE: func(c *cobra.Command, args []string) error {
all, err := kdb.RecentEntries(limit * 10) if page < 1 || limit < 1 {
return fmt.Errorf("page and limit must be >= 1 (got page=%d limit=%d)", page, limit)
}
all, err := kdb.RecentEntries(page * limit)
if err != nil { if err != nil {
return err return err
} }
entries := paginate(all, page, limit) entries := paginate(all, page, limit)
if len(entries) == 0 { if len(entries) == 0 {
fmt.Println("No entries yet. Run `knox watch` to start ingesting.") if page > 1 {
fmt.Printf("No more entries on page %d (page %d of %d).\n", page, page, (len(all)+limit-1)/limit)
} else {
fmt.Println("No entries yet. Run `knox watch` to start ingesting.")
}
return nil return nil
} }
fmt.Printf("Recent %d entries (page %d, %d per page):\n\n", len(entries), page, limit) fmt.Printf("Recent %d entries (page %d, %d per page):\n\n", len(entries), page, limit)
+1 -1
View File
@@ -5,9 +5,9 @@ import (
"regexp" "regexp"
"strings" "strings"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/index" "github.com/david/knox/internal/index"
"github.com/spf13/cobra"
) )
func NewTopicsCmd(kdb *db.KnoxDB) *cobra.Command { func NewTopicsCmd(kdb *db.KnoxDB) *cobra.Command {
+7 -2
View File
@@ -1,12 +1,14 @@
package cmd package cmd
import ( import (
"fmt"
"io" "io"
"log" "log"
"os"
"github.com/spf13/cobra"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
"github.com/david/knox/internal/watch" "github.com/david/knox/internal/watch"
"github.com/spf13/cobra"
) )
func NewWatchCmd(kdb *db.KnoxDB) *cobra.Command { func NewWatchCmd(kdb *db.KnoxDB) *cobra.Command {
@@ -25,7 +27,10 @@ func NewWatchCmd(kdb *db.KnoxDB) *cobra.Command {
log.Printf("[knox] starting watcher, scanning %d directories", len(dirs)) log.Printf("[knox] starting watcher, scanning %d directories", len(dirs))
w := watch.New(kdb, dirs) w := watch.New(kdb, dirs)
if err := w.Start(); err != nil { if err := w.Start(); err != nil {
log.Fatalf("watcher error: %v", err) // Not log.Fatalf: --quiet redirects the logger to io.Discard,
// so a fatal exit must reach the user on stderr directly.
fmt.Fprintf(os.Stderr, "watcher error: %v\n", err)
os.Exit(1)
} }
}, },
} }
+164 -73
View File
@@ -22,6 +22,12 @@ type KnoxDB struct {
db *sql.DB db *sql.DB
nodeID string nodeID string
clock *hlc.Clock clock *hlc.Clock
// writeMu serializes the check-then-insert dedup in RecordObservation within
// this process. Cross-process serialization comes from _txlock=immediate (the
// write lock is taken at BEGIN, before the dedup read) plus a single
// connection per pool.
writeMu sync.Mutex
} }
type Observation struct { type Observation struct {
@@ -93,10 +99,13 @@ func Open(path string) (*KnoxDB, error) {
return nil, fmt.Errorf("create db dir: %w", err) return nil, fmt.Errorf("create db dir: %w", err)
} }
db, err := sql.Open("sqlite", path+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)") db, err := sql.Open("sqlite", path+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)&_txlock=immediate")
if err != nil { if err != nil {
return nil, fmt.Errorf("open db: %w", err) return nil, fmt.Errorf("open db: %w", err)
} }
// One connection per process: WAL has a single writer; serializing on one
// connection avoids pool contention surfacing as busy_timeout errors.
db.SetMaxOpenConns(1)
if _, err := db.Exec(Schema); err != nil { if _, err := db.Exec(Schema); err != nil {
return nil, fmt.Errorf("init schema: %w", err) return nil, fmt.Errorf("init schema: %w", err)
@@ -135,6 +144,13 @@ func Open(path string) (*KnoxDB, error) {
if _, err := db.Exec("UPDATE observations SET hcl=id WHERE hcl IS NULL"); err != nil { if _, err := db.Exec("UPDATE observations SET hcl=id WHERE hcl IS NULL"); err != nil {
return nil, fmt.Errorf("backfill hcl: %w", err) return nil, fmt.Errorf("backfill hcl: %w", err)
} }
// Resume this node's HLC from its persisted max: a restart with a regressed
// wall clock must not reissue already-persisted values (see hlc.SeekTo).
var maxHCL int64
if err := db.QueryRow(`SELECT COALESCE(MAX(hcl), 0) FROM observations WHERE node_id=?`, nodeID).Scan(&maxHCL); err != nil {
return nil, fmt.Errorf("seed hlc: %w", err)
}
kdb.clock.SeekTo(maxHCL)
// Locator uniqueness: (node_id, hcl) is the merge key for gossip; a given // Locator uniqueness: (node_id, hcl) is the merge key for gossip; a given
// node's HCL is strictly monotonic so this never throws a false conflict. // node's HCL is strictly monotonic so this never throws a false conflict.
@@ -204,6 +220,9 @@ func (k *KnoxDB) ObservationEntryEstimate() int {
// Idempotent: if the latest observation for this fingerprint has an identical // Idempotent: if the latest observation for this fingerprint has an identical
// content signature, nothing is recorded — re-ingesting unchanged content is a no-op. // content signature, nothing is recorded — re-ingesting unchanged content is a no-op.
func (k *KnoxDB) RecordObservation(o ObservationRecord) (obsID int64, isNew bool, err error) { func (k *KnoxDB) RecordObservation(o ObservationRecord) (obsID int64, isNew bool, err error) {
k.writeMu.Lock()
defer k.writeMu.Unlock()
tx, err := k.db.Begin() tx, err := k.db.Begin()
if err != nil { if err != nil {
return 0, false, fmt.Errorf("begin tx: %w", err) return 0, false, fmt.Errorf("begin tx: %w", err)
@@ -289,18 +308,21 @@ func (k *KnoxDB) RecordObservation(o ObservationRecord) (obsID int64, isNew bool
`INSERT OR IGNORE INTO thread_observations (thread_id, observation_id, relevance) VALUES (?, ?, 'auto')`, `INSERT OR IGNORE INTO thread_observations (thread_id, observation_id, relevance) VALUES (?, ?, 'auto')`,
goldenID, obsID, goldenID, obsID,
) )
linked := err == nil && res != nil if err != nil {
if linked { return 0, false, fmt.Errorf("auto-link golden thread: %w", err)
if n, _ := res.RowsAffected(); n == 0 { }
linked = false affected, _ := res.RowsAffected()
} linked := affected > 0
if _, err := tx.Exec(`UPDATE threads SET updated_at=datetime('now') WHERE id=?`, goldenID); err != nil {
return 0, false, fmt.Errorf("bump golden thread updated_at: %w", err)
} }
tx.Exec(`UPDATE threads SET updated_at=datetime('now') WHERE id=?`, goldenID)
// Temporal cross-reference: when a high-signal non-browser event links, // Temporal cross-reference: when a high-signal non-browser event links,
// nearby browser history (±1h) becomes thread context too. // nearby browser history (±1h) becomes thread context too.
if linked && o.SourceID != "browser-history" && o.CreatedAt != "" { if linked && o.SourceID != "browser-history" && o.CreatedAt != "" {
k.linkTemporalNeighbors(tx, goldenID, o.CreatedAt, keywords, bigrams) if err := k.linkTemporalNeighbors(tx, goldenID, o.CreatedAt, keywords, bigrams); err != nil {
return 0, false, fmt.Errorf("link temporal neighbors: %w", err)
}
} }
} }
} }
@@ -311,7 +333,7 @@ func (k *KnoxDB) RecordObservation(o ObservationRecord) (obsID int64, isNew bool
// linkTemporalNeighbors links browser-history observations within ±1h of a // linkTemporalNeighbors links browser-history observations within ±1h of a
// signal timestamp to the thread. A weak keyword check still applies so // signal timestamp to the thread. A weak keyword check still applies so
// unrelated browsing (lunch news) doesn't pollute the thread. // unrelated browsing (lunch news) doesn't pollute the thread.
func (k *KnoxDB) linkTemporalNeighbors(tx *sql.Tx, threadID int64, createdAt string, keywords []string, bigrams [][2]string) { func (k *KnoxDB) linkTemporalNeighbors(tx *sql.Tx, threadID int64, createdAt string, keywords []string, bigrams [][2]string) error {
rows, err := tx.Query( rows, err := tx.Query(
`SELECT o.id, COALESCE(o.title,''), COALESCE(o.summary,'') `SELECT o.id, COALESCE(o.title,''), COALESCE(o.summary,'')
FROM observations o FROM observations o
@@ -322,7 +344,7 @@ func (k *KnoxDB) linkTemporalNeighbors(tx *sql.Tx, threadID int64, createdAt str
createdAt, createdAt,
) )
if err != nil { if err != nil {
return return err
} }
defer rows.Close() defer rows.Close()
@@ -330,15 +352,18 @@ func (k *KnoxDB) linkTemporalNeighbors(tx *sql.Tx, threadID int64, createdAt str
var obsID int64 var obsID int64
var title, summary string var title, summary string
if err := rows.Scan(&obsID, &title, &summary); err != nil { if err := rows.Scan(&obsID, &title, &summary); err != nil {
continue return err
} }
if isRelevant(title, summary, "", keywords, bigrams) || weakRelevant(title, summary, keywords) { if isRelevant(title, summary, "", keywords, bigrams) || weakRelevant(title, summary, keywords) {
tx.Exec( if _, err := tx.Exec(
`INSERT OR IGNORE INTO thread_observations (thread_id, observation_id, relevance) VALUES (?, ?, 'temporal')`, `INSERT OR IGNORE INTO thread_observations (thread_id, observation_id, relevance) VALUES (?, ?, 'temporal')`,
threadID, obsID, threadID, obsID,
) ); err != nil {
return err
}
} }
} }
return rows.Err()
} }
// weakRelevant is a lower bar for temporal neighbors: any single keyword match. // weakRelevant is a lower bar for temporal neighbors: any single keyword match.
@@ -428,23 +453,39 @@ func (k *KnoxDB) RebuildEntriesFromObservations() (before, after int, err error)
if _, err := tx.Exec("DELETE FROM entries"); err != nil { if _, err := tx.Exec("DELETE FROM entries"); err != nil {
return 0, 0, err return 0, 0, err
} }
// Content fields come from the latest observation per fingerprint, folded in
// a total order (hcl DESC, node_id DESC, id DESC) — deterministic for any two
// nodes holding identical logs, per the gossip spec (bare GROUP BY columns
// would take an arbitrary row, merge-order dependent).
_, err = tx.Exec( _, err = tx.Exec(
`INSERT INTO entries (fingerprint, source_id, source_path, project, content_type, title, summary, first_seen, last_seen, created_at, ref_count, last_confidence) `INSERT INTO entries (fingerprint, source_id, source_path, project, content_type, title, summary, first_seen, last_seen, created_at, ref_count, last_confidence)
SELECT SELECT
o.fingerprint, latest.fingerprint,
o.source_id, latest.source_id,
o.source_path, latest.source_path,
o.project, latest.project,
o.content_type, latest.content_type,
o.title, latest.title,
o.summary, latest.summary,
MIN(o.collected_at), agg.first_seen,
MAX(o.collected_at), agg.last_seen,
COALESCE(NULLIF(MAX(o.created_at),''), MIN(o.collected_at)), agg.created_at,
COUNT(*), agg.ref_count,
MAX(o.confidence) agg.max_confidence
FROM observations o FROM (
GROUP BY o.fingerprint`, SELECT *, ROW_NUMBER() OVER (PARTITION BY fingerprint ORDER BY hcl DESC, node_id DESC, id DESC) AS rn
FROM observations
) latest
JOIN (
SELECT fingerprint,
MIN(collected_at) AS first_seen,
MAX(collected_at) AS last_seen,
COALESCE(NULLIF(MAX(created_at),''), MIN(collected_at)) AS created_at,
COUNT(*) AS ref_count,
MAX(confidence) AS max_confidence
FROM observations GROUP BY fingerprint
) agg ON agg.fingerprint = latest.fingerprint
WHERE latest.rn = 1`,
) )
if err != nil { if err != nil {
return 0, 0, fmt.Errorf("rebuild entries: %w", err) return 0, 0, fmt.Errorf("rebuild entries: %w", err)
@@ -585,6 +626,9 @@ func scanEntries(rows *sql.Rows) ([]Entry, error) {
} }
entries = append(entries, e) entries = append(entries, e)
} }
if err := rows.Err(); err != nil {
return nil, err
}
return entries, nil return entries, nil
} }
@@ -636,38 +680,50 @@ func (k *KnoxDB) MarkSessionIndexed(sessionID string) error {
func (k *KnoxDB) Stats() (map[string]any, error) { func (k *KnoxDB) Stats() (map[string]any, error) {
stats := make(map[string]any) stats := make(map[string]any)
var v int // A broken DB reports an error, not zeros: every aggregate is checked.
var s string countInt := func(query string) (int, error) {
var v int
if err := k.db.QueryRow(query).Scan(&v); err != nil {
return 0, err
}
return v, nil
}
ints := []struct {
key string
query string
}{
{"total_entries", "SELECT COUNT(*) FROM entries"},
{"entries_last_24h", "SELECT COUNT(*) FROM entries WHERE last_seen > datetime('now', '-1 day')"},
{"total_projects", "SELECT COUNT(DISTINCT project) FROM entries"},
{"total_observations", "SELECT COUNT(*) FROM observations"},
{"observations_last_24h", "SELECT COUNT(*) FROM observations WHERE collected_at > datetime('now', '-1 day')"},
{"total_sessions", "SELECT COUNT(*) FROM sessions"},
{"pending_reflections", "SELECT COUNT(*) FROM sessions WHERE indexed=0"},
}
for _, it := range ints {
v, err := countInt(it.query)
if err != nil {
return nil, fmt.Errorf("stats %s: %w", it.key, err)
}
stats[it.key] = v
}
k.db.QueryRow("SELECT COUNT(*) FROM entries").Scan(&v) var earliest string
stats["total_entries"] = v if err := k.db.QueryRow("SELECT COALESCE(MIN(collected_at),'') FROM observations").Scan(&earliest); err != nil {
v = 0 return nil, fmt.Errorf("stats earliest_observation: %w", err)
k.db.QueryRow("SELECT COUNT(*) FROM entries WHERE last_seen > datetime('now', '-1 day')").Scan(&v) }
stats["entries_last_24h"] = v stats["earliest_observation"] = earliest
v = 0
k.db.QueryRow("SELECT COUNT(DISTINCT project) FROM entries").Scan(&v)
stats["total_projects"] = v
v = 0
k.db.QueryRow("SELECT COUNT(*) FROM observations").Scan(&v)
stats["total_observations"] = v
v = 0
k.db.QueryRow("SELECT COUNT(*) FROM observations WHERE collected_at > datetime('now', '-1 day')").Scan(&v)
stats["observations_last_24h"] = v
v = 0
k.db.QueryRow("SELECT COALESCE(MIN(collected_at),'') FROM observations").Scan(&s)
stats["earliest_observation"] = s
s = ""
k.db.QueryRow("SELECT COUNT(*) FROM sessions").Scan(&v)
stats["total_sessions"] = v
v = 0
k.db.QueryRow("SELECT COUNT(*) FROM sessions WHERE indexed=0").Scan(&v)
stats["pending_reflections"] = v
v = 0
goldenID, _ := k.GoldenThreadID() goldenID, err := k.GoldenThreadID()
if err != nil {
return nil, fmt.Errorf("stats golden_thread: %w", err)
}
if goldenID > 0 { if goldenID > 0 {
stats["golden_thread_id"] = goldenID stats["golden_thread_id"] = goldenID
t, _ := k.GetThread(goldenID) t, err := k.GetThread(goldenID)
if err != nil {
return nil, fmt.Errorf("stats golden_thread: %w", err)
}
if t != nil { if t != nil {
stats["golden_thread"] = t.Title stats["golden_thread"] = t.Title
} }
@@ -970,20 +1026,26 @@ func (k *KnoxDB) CreateThread(title, motivation, priority, tags, provenance stri
} }
func (k *KnoxDB) CreateThreadCluster(title, motivation, priority, tags, provenance, clusterKey string) (int64, bool, error) { func (k *KnoxDB) CreateThreadCluster(title, motivation, priority, tags, provenance, clusterKey string) (int64, bool, error) {
if clusterKey != "" { // INSERT OR IGNORE makes creation atomic under the partial unique index:
if existing := k.ThreadByClusterKey(clusterKey); existing != 0 { // concurrent creators (thread ticker + post-pull reconcile, or two gossip
return existing, false, nil // nodes converging on the same cluster) lose the race into a no-op instead
} // of a UNIQUE error, then pick up the winner's id.
}
res, err := k.db.Exec( res, err := k.db.Exec(
`INSERT INTO threads (title, motivation, priority, tags, provenance, cluster_key) VALUES (?, ?, ?, ?, ?, ?)`, `INSERT OR IGNORE INTO threads (title, motivation, priority, tags, provenance, cluster_key) VALUES (?, ?, ?, ?, ?, ?)`,
title, motivation, priority, tags, provenance, clusterKey, title, motivation, priority, tags, provenance, clusterKey,
) )
if err != nil { if err != nil {
return 0, false, err return 0, false, err
} }
id, _ := res.LastInsertId() if n, _ := res.RowsAffected(); n > 0 {
return id, true, nil id, _ := res.LastInsertId()
return id, true, nil
}
// Ignored insert: a thread with this cluster key already exists.
if existing := k.ThreadByClusterKey(clusterKey); existing != 0 {
return existing, false, nil
}
return 0, false, fmt.Errorf("insert ignored but no thread with cluster key %q", clusterKey)
} }
// ThreadByClusterKey returns the id of the thread auto-created for a given // ThreadByClusterKey returns the id of the thread auto-created for a given
@@ -1088,15 +1150,21 @@ func (k *KnoxDB) UpdateThread(id int64, title, motivation, priority, tags string
} }
func (k *KnoxDB) LinkObservationToThread(threadID, observationID int64, relevance string) error { func (k *KnoxDB) LinkObservationToThread(threadID, observationID int64, relevance string) error {
_, err := k.db.Exec( tx, err := k.db.Begin()
`INSERT OR IGNORE INTO thread_observations (thread_id, observation_id, relevance) VALUES (?, ?, ?)`,
threadID, observationID, relevance,
)
if err != nil { if err != nil {
return err return err
} }
_, err = k.db.Exec(`UPDATE threads SET updated_at=datetime('now') WHERE id=?`, threadID) defer tx.Rollback()
return err if _, err := tx.Exec(
`INSERT OR IGNORE INTO thread_observations (thread_id, observation_id, relevance) VALUES (?, ?, ?)`,
threadID, observationID, relevance,
); err != nil {
return err
}
if _, err := tx.Exec(`UPDATE threads SET updated_at=datetime('now') WHERE id=?`, threadID); err != nil {
return err
}
return tx.Commit()
} }
// AutoLinkThreadObservations resolves each entry fingerprint to its most recent // AutoLinkThreadObservations resolves each entry fingerprint to its most recent
@@ -1111,6 +1179,7 @@ func (k *KnoxDB) AutoLinkThreadObservations(threadID int64, fingerprints []strin
defer tx.Rollback() defer tx.Rollback()
linked := 0 linked := 0
var firstErr error
for _, fp := range fingerprints { for _, fp := range fingerprints {
if fp == "" { if fp == "" {
continue continue
@@ -1119,6 +1188,11 @@ func (k *KnoxDB) AutoLinkThreadObservations(threadID int64, fingerprints []strin
if err := tx.QueryRow( if err := tx.QueryRow(
`SELECT id FROM observations WHERE fingerprint=? ORDER BY hcl DESC, id DESC LIMIT 1`, fp, `SELECT id FROM observations WHERE fingerprint=? ORDER BY hcl DESC, id DESC LIMIT 1`, fp,
).Scan(&obsID); err != nil { ).Scan(&obsID); err != nil {
// A fingerprint without observations is expected (provenance links);
// a real query failure is not — remember the first one.
if err != sql.ErrNoRows && firstErr == nil {
firstErr = err
}
continue continue
} }
res, err := tx.Exec( res, err := tx.Exec(
@@ -1126,6 +1200,9 @@ func (k *KnoxDB) AutoLinkThreadObservations(threadID int64, fingerprints []strin
threadID, obsID, threadID, obsID,
) )
if err != nil { if err != nil {
if firstErr == nil {
firstErr = err
}
continue continue
} }
if n, _ := res.RowsAffected(); n > 0 { if n, _ := res.RowsAffected(); n > 0 {
@@ -1138,7 +1215,10 @@ func (k *KnoxDB) AutoLinkThreadObservations(threadID int64, fingerprints []strin
return linked, err return linked, err
} }
} }
return linked, tx.Commit() if err := tx.Commit(); err != nil {
return linked, err
}
return linked, firstErr
} }
// ActiveThreadByKeyword returns the most recently updated active thread whose // ActiveThreadByKeyword returns the most recently updated active thread whose
@@ -1207,15 +1287,26 @@ func (k *KnoxDB) ThreadObservations(threadID int64) ([]Observation, error) {
} }
func (k *KnoxDB) AddThreadNote(threadID int64, note string) (int64, error) { func (k *KnoxDB) AddThreadNote(threadID int64, note string) (int64, error) {
res, err := k.db.Exec( tx, err := k.db.Begin()
if err != nil {
return 0, err
}
defer tx.Rollback()
res, err := tx.Exec(
`INSERT INTO thread_notes (thread_id, note) VALUES (?, ?)`, `INSERT INTO thread_notes (thread_id, note) VALUES (?, ?)`,
threadID, note, threadID, note,
) )
if err != nil { if err != nil {
return 0, err return 0, err
} }
_, err = k.db.Exec(`UPDATE threads SET updated_at=datetime('now') WHERE id=?`, threadID) if _, err := tx.Exec(`UPDATE threads SET updated_at=datetime('now') WHERE id=?`, threadID); err != nil {
return res.LastInsertId() return 0, err
}
id, err := res.LastInsertId()
if err != nil {
return 0, err
}
return id, tx.Commit()
} }
func (k *KnoxDB) ThreadNotes(threadID int64) ([]struct { func (k *KnoxDB) ThreadNotes(threadID int64) ([]struct {
+242
View File
@@ -0,0 +1,242 @@
package db
import (
"path/filepath"
"sync"
"testing"
)
func tmpDB(t *testing.T) *KnoxDB {
t.Helper()
kdb, err := Open(filepath.Join(t.TempDir(), "index.db"))
if err != nil {
t.Fatalf("open db: %v", err)
}
t.Cleanup(func() { kdb.Close() })
return kdb
}
func record(t *testing.T, k *KnoxDB, fp, title string) {
t.Helper()
_, _, err := k.RecordObservation(ObservationRecord{
Fingerprint: fp,
SourceID: "test",
SourcePath: fp,
Project: "itest",
ContentType: "test",
Title: title,
Summary: "summary",
CreatedAt: "2026-08-29T00:00:00Z",
LineEnd: 0,
Confidence: 0.9,
IngesterVersion: "test/v1",
})
if err != nil {
t.Errorf("record %s: %v", fp, err)
}
}
func obsCount(t *testing.T, k *KnoxDB) int {
t.Helper()
stats, err := k.Stats()
if err != nil {
t.Fatalf("stats: %v", err)
}
n, _ := stats["total_observations"].(int)
return n
}
// TestRecordObservationConcurrentDedup: N goroutines ingesting identical
// content must produce exactly one observation row. This exercises the
// check-then-insert dedup under the writeMu + BEGIN IMMEDIATE serialization.
func TestRecordObservationConcurrentDedup(t *testing.T) {
k := tmpDB(t)
var wg sync.WaitGroup
for i := 0; i < 20; i++ {
wg.Add(1)
go func() {
defer wg.Done()
record(t, k, "fp:concurrent", "same title")
}()
}
wg.Wait()
if n := obsCount(t, k); n != 1 {
t.Fatalf("expected exactly 1 observation after concurrent identical ingests, got %d", n)
}
}
// TestRecordObservationDistinctFingerprints: different content must never be
// deduped away (the constraint is per-fingerprint signature, not global).
func TestRecordObservationDistinctFingerprints(t *testing.T) {
k := tmpDB(t)
record(t, k, "fp:a", "alpha")
record(t, k, "fp:b", "beta")
record(t, k, "fp:a", "alpha changed")
if n := obsCount(t, k); n != 3 {
t.Fatalf("expected 3 observations, got %d", n)
}
}
// TestPushObservationsValidation: forged/malformed rows are rejected without
// error — self node_id (poisoning vector), empty node_id, negative HCL.
func TestPushObservationsValidation(t *testing.T) {
k := tmpDB(t)
foreign := GossipObservation{
NodeID: "0123456789abcdef0123456789abcdef", HCL: 42,
Fingerprint: "fp:foreign", SourceID: "test", Title: "t", Summary: "s",
CollectedAt: "2026-08-29T00:00:00Z",
}
rows := []GossipObservation{
foreign,
{NodeID: k.NodeID(), HCL: 100, Fingerprint: "fp:self"}, // spoof poisoning attempt
{NodeID: "", HCL: 1, Fingerprint: "fp:empty"},
{NodeID: "other", HCL: -5, Fingerprint: "fp:neg"},
}
n, err := k.PushObservations(rows)
if err != nil {
t.Fatalf("push: %v", err)
}
if n != 1 {
t.Fatalf("expected exactly the one valid row inserted, got %d", n)
}
if got := obsCount(t, k); got != 1 {
t.Fatalf("expected 1 observation in the log, got %d", got)
}
}
// TestOpenReopenHCLMonotonicAcrossRestart: reopening a DB must resume the HLC
// from its persisted max (clock seeding), keep the node identity, and order the
// new observation above every previous one.
func TestOpenReopenHCLMonotonicAcrossRestart(t *testing.T) {
path := filepath.Join(t.TempDir(), "index.db")
k1, err := Open(path)
if err != nil {
t.Fatalf("open: %v", err)
}
record(t, k1, "fp:r1", "one")
record(t, k1, "fp:r2", "two")
record(t, k1, "fp:r3", "three")
nodeID1 := k1.NodeID()
maxBefore := int64(0)
rows, err := k1.ObservationsAfter(nodeID1, 0, 100)
if err != nil {
t.Fatalf("obs after: %v", err)
}
for _, r := range rows {
if r.HCL > maxBefore {
maxBefore = r.HCL
}
}
if err := k1.Close(); err != nil {
t.Fatalf("close: %v", err)
}
k2, err := Open(path)
if err != nil {
t.Fatalf("reopen: %v", err)
}
defer k2.Close()
if k2.NodeID() != nodeID1 {
t.Errorf("node id changed across reopen: %q -> %q", nodeID1, k2.NodeID())
}
record(t, k2, "fp:r4", "four")
after, err := k2.ObservationsAfter(nodeID1, maxBefore, 100)
if err != nil {
t.Fatalf("obs after (reopened): %v", err)
}
if len(after) != 1 {
t.Fatalf("expected exactly the new observation above the pre-restart max, got %d rows", len(after))
}
if after[0].Fingerprint != "fp:r4" {
t.Errorf("unexpected row above max: %s", after[0].Fingerprint)
}
}
// TestCreateThreadClusterConcurrentIdempotent: two concurrent creators of the
// same cluster key must converge on one thread (one created=true, one false)
// with no UNIQUE constraint error — the INSERT OR IGNORE path.
func TestCreateThreadClusterConcurrentIdempotent(t *testing.T) {
k := tmpDB(t)
const key = "cluster:race"
var wg sync.WaitGroup
ids := make([]int64, 2)
created := make([]bool, 2)
errs := make([]error, 2)
for i := 0; i < 2; i++ {
wg.Add(1)
go func(i int) {
defer wg.Done()
ids[i], created[i], errs[i] = k.CreateThreadCluster("Race thread", "m", "medium", "[]", `{}`, key)
}(i)
}
wg.Wait()
for i := range errs {
if errs[i] != nil {
t.Fatalf("creator %d: %v", i, errs[i])
}
}
if ids[0] != ids[1] {
t.Errorf("concurrent creators got different ids: %d vs %d", ids[0], ids[1])
}
if created[0] == created[1] {
t.Errorf("exactly one creator should report created=true, got %v %v", created[0], created[1])
}
threads, err := k.ListThreads("")
if err != nil {
t.Fatal(err)
}
n := 0
for _, th := range threads {
if th.Title == "Race thread" {
n++
}
}
if n != 1 {
t.Errorf("want exactly 1 race thread, got %d", n)
}
}
// TestRebuildEntriesDeterministicFold: two nodes holding the same observations
// in different merge orders must rebuild identical derived state — content
// fields come from the highest-HCL observation per fingerprint (spec 6.2),
// never from an arbitrary (merge-order-dependent) row.
func TestRebuildEntriesDeterministicFold(t *testing.T) {
const foreign = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
rows := []GossipObservation{
{NodeID: foreign, HCL: 10, Fingerprint: "fp:fold", SourceID: "test", Title: "old title", Summary: "old", CollectedAt: "2026-08-29T00:00:00Z", Confidence: 0.5},
{NodeID: foreign, HCL: 20, Fingerprint: "fp:fold", SourceID: "test", Title: "new title", Summary: "new", CollectedAt: "2026-08-29T01:00:00Z", Confidence: 0.9},
}
results := make([]Entry, 2)
for i, order := range [][]GossipObservation{{rows[0], rows[1]}, {rows[1], rows[0]}} {
k := tmpDB(t)
if _, err := k.PushObservations(order); err != nil {
t.Fatalf("push: %v", err)
}
if _, _, err := k.RebuildEntriesFromObservations(); err != nil {
t.Fatalf("rebuild: %v", err)
}
e, err := k.FindEntry("fp:fold")
if err != nil || e == nil {
t.Fatalf("find entry: %v", err)
}
results[i] = *e
}
if results[0].Title != "new title" {
t.Errorf("order [old,new]: title = %q, want %q", results[0].Title, "new title")
}
if results[1].Title != "new title" {
t.Errorf("order [new,old]: title = %q, want %q", results[1].Title, "new title")
}
if results[0] != results[1] {
t.Errorf("derived state diverged between merge orders:\n%+v\n%+v", results[0], results[1])
}
if results[0].RefCount != 2 {
t.Errorf("ref_count = %d, want 2", results[0].RefCount)
}
if results[0].FirstSeen != "2026-08-29T00:00:00Z" || results[0].LastSeen != "2026-08-29T01:00:00Z" {
t.Errorf("first/last seen = %q/%q", results[0].FirstSeen, results[0].LastSeen)
}
}
+14
View File
@@ -3,6 +3,7 @@ package db
import ( import (
"database/sql" "database/sql"
"fmt" "fmt"
"log"
) )
// GossipObservation is the serializable wire form of an observation exchanged // GossipObservation is the serializable wire form of an observation exchanged
@@ -40,6 +41,19 @@ func (k *KnoxDB) PushObservations(rows []GossipObservation) (int, error) {
inserted := 0 inserted := 0
for _, o := range rows { for _, o := range rows {
// Reject malformed or forged rows. Nodes only ever push their own
// observations, so a row claiming this node's id cannot be legitimate:
// accepting it would let a peer poison our knowledge vector (a forged
// max-HCL makes peers believe they have our whole history and stop
// pulling). Empty node ids and negative HCLs are likewise never produced
// by a real node.
if o.NodeID == "" || o.HCL < 0 {
continue
}
if o.NodeID == k.nodeID {
log.Printf("[gossip] dropped pushed row claiming local node_id (spoof?)")
continue
}
res, err := tx.Exec( res, err := tx.Exec(
`INSERT OR IGNORE INTO observations `INSERT OR IGNORE INTO observations
(fingerprint, source_id, source_path, project, content_type, title, summary, (fingerprint, source_id, source_path, project, content_type, title, summary,
+1 -1
View File
@@ -11,7 +11,7 @@ type MetricsSnapshot struct {
PendingReflections int PendingReflections int
ThreadsByStatus map[string]int ThreadsByStatus map[string]int
Peers int Peers int
ByOriginNode map[string]int // node_id → observation count ByOriginNode map[string]int // node_id → observation count
KnowledgeVector map[string]int64 // node_id → max hcl KnowledgeVector map[string]int64 // node_id → max hcl
EarliestObservation string EarliestObservation string
} }
+17
View File
@@ -27,6 +27,23 @@ type Clock struct {
func New() *Clock { return &Clock{} } func New() *Clock { return &Clock{} }
// SeekTo adopts the given packed HLC value when it is ahead of the clock's current
// position, so the next Now is still strictly increasing. Used to resume a node's
// clock from its persisted MAX(hcl) at startup — without it, a restart with a
// regressed wall clock would reissue already-used values and break the
// monotonicity the (node_id, hcl) locator uniqueness and gossip cursors rely on.
func (c *Clock) SeekTo(v int64) {
c.mu.Lock()
defer c.mu.Unlock()
wall := v >> wallShift
seq := v & seqMask
if wall > c.wallMS || (wall == c.wallMS && seq > c.seq) {
c.wallMS = wall
c.seq = seq
}
}
// Now returns the next monotonic HLC value and the wall-clock time embedded in // Now returns the next monotonic HLC value and the wall-clock time embedded in
// it. The returned time is the HLC's wall component — never ahead of the local // it. The returned time is the HLC's wall component — never ahead of the local
// clock beyond the current call and never rewinding across calls. // clock beyond the current call and never rewinding across calls.
+68
View File
@@ -0,0 +1,68 @@
package hlc
import "testing"
// TestSeekToFutureValueKeepsMonotonic: after resuming from a persisted value
// ahead of the wall clock (clock regression), every subsequent Now must still
// be strictly greater than the resumed position.
func TestSeekToFutureValueKeepsMonotonic(t *testing.T) {
c := New()
resumed := int64(1) << 62 // packed value far ahead of any real wall clock
c.SeekTo(resumed)
prev, _ := c.Now()
if prev <= resumed {
t.Fatalf("first Now after SeekTo = %d, want > resumed %d", prev, resumed)
}
for i := 0; i < 100; i++ {
next, _ := c.Now()
if next <= prev {
t.Fatalf("HLC regressed: %d then %d", prev, next)
}
prev = next
}
}
// TestSeekToLowerValueIgnored: resuming from a value behind the current clock
// (or a fresh 0-padded DB) must not rewind it.
func TestSeekToLowerValueIgnored(t *testing.T) {
c := New()
first, _ := c.Now()
c.SeekTo(0)
second, _ := c.Now()
if second <= first {
t.Fatalf("SeekTo(0) rewound the clock: %d then %d", first, second)
}
// Seek to exactly the last emitted value: the next value must exceed it.
c.SeekTo(second)
third, _ := c.Now()
if third <= second {
t.Fatalf("SeekTo(last) did not preserve monotonicity: %d then %d", second, third)
}
}
// TestSeekToAcrossRestart mirrors Open's reopen path: a fresh clock resumed
// from the persisted max keeps issuing strictly increasing values.
func TestSeekToAcrossRestart(t *testing.T) {
c1 := New()
var last int64
for i := 0; i < 50; i++ {
last, _ = c1.Now()
}
c2 := New() // fresh process clock
c2.SeekTo(last)
prev, _ := c2.Now()
if prev <= last {
t.Fatalf("reopened clock reissued a value: %d <= %d", prev, last)
}
for i := 0; i < 50; i++ {
next, _ := c2.Now()
if next <= prev {
t.Fatalf("reopened clock regressed: %d then %d", prev, next)
}
prev = next
}
}
+21 -21
View File
@@ -19,23 +19,23 @@ import (
// scrape-time gauges derived from the database (refreshed on each scrape). // scrape-time gauges derived from the database (refreshed on each scrape).
type Metrics struct { type Metrics struct {
// Gossip counters (live). // Gossip counters (live).
pullsTotal prometheus.Counter pullsTotal prometheus.Counter
pushesTotal prometheus.Counter pushesTotal prometheus.Counter
obsPulledTotal prometheus.Counter obsPulledTotal prometheus.Counter
obsPushedTotal prometheus.Counter obsPushedTotal prometheus.Counter
errorsTotal prometheus.Counter errorsTotal prometheus.Counter
// Snapshot gauges (updated per scrape). // Snapshot gauges (updated per scrape).
observationsGauge *prometheus.GaugeVec observationsGauge *prometheus.GaugeVec
entriesGauge prometheus.Gauge entriesGauge prometheus.Gauge
projectsGauge prometheus.Gauge projectsGauge prometheus.Gauge
sessionsGauge prometheus.Gauge sessionsGauge prometheus.Gauge
pendingReflections prometheus.Gauge pendingReflections prometheus.Gauge
peersGauge prometheus.Gauge peersGauge prometheus.Gauge
threadsByStatus *prometheus.GaugeVec threadsByStatus *prometheus.GaugeVec
byOriginNode *prometheus.GaugeVec byOriginNode *prometheus.GaugeVec
knowledgeVector *prometheus.GaugeVec knowledgeVector *prometheus.GaugeVec
observationsLast24h prometheus.Gauge observationsLast24h prometheus.Gauge
registry *prometheus.Registry registry *prometheus.Registry
kdb *db.KnoxDB kdb *db.KnoxDB
@@ -59,14 +59,14 @@ func New(kdb *db.KnoxDB, nodeName string) *Metrics {
m.obsPushedTotal = newCounter(reg, "knox_gossip_observations_pushed_total", "Observations sent to peers.") m.obsPushedTotal = newCounter(reg, "knox_gossip_observations_pushed_total", "Observations sent to peers.")
m.errorsTotal = newCounter(reg, "knox_gossip_errors_total", "Gossip errors (ping/pull/push failures).") m.errorsTotal = newCounter(reg, "knox_gossip_errors_total", "Gossip errors (ping/pull/push failures).")
m.observationsGauge = newGaugeVec(reg, "knox_observations_total", "Observation log size.", "source_id") m.observationsGauge = newGaugeVec(reg, "knox_observations", "Observation log size.", "source_id")
m.observationsLast24h = newGauge(reg, "knox_observations_last_24h", "Observations collected in the last 24h.") m.observationsLast24h = newGauge(reg, "knox_observations_last_24h", "Observations collected in the last 24h.")
m.entriesGauge = newGauge(reg, "knox_entries_total", "Materialized entry cache size.") m.entriesGauge = newGauge(reg, "knox_entries", "Materialized entry cache size.")
m.projectsGauge = newGauge(reg, "knox_projects_total", "Distinct projects in the entry cache.") m.projectsGauge = newGauge(reg, "knox_projects", "Distinct projects in the entry cache.")
m.sessionsGauge = newGauge(reg, "knox_sessions_total", "Sessions tracked.") m.sessionsGauge = newGauge(reg, "knox_sessions", "Sessions tracked.")
m.pendingReflections = newGauge(reg, "knox_pending_reflections", "Sessions awaiting reflection.") m.pendingReflections = newGauge(reg, "knox_pending_reflections", "Sessions awaiting reflection.")
m.peersGauge = newGauge(reg, "knox_peers_total", "Known peer nodes.") m.peersGauge = newGauge(reg, "knox_peers", "Known peer nodes.")
m.threadsByStatus = newGaugeVec(reg, "knox_threads_total", "Threads by status.", "status") m.threadsByStatus = newGaugeVec(reg, "knox_threads", "Threads by status.", "status")
m.byOriginNode = newGaugeVec(reg, "knox_observations_by_node", "Observations per originating node.", "node_id") m.byOriginNode = newGaugeVec(reg, "knox_observations_by_node", "Observations per originating node.", "node_id")
m.knowledgeVector = newGaugeVec(reg, "knox_knowledge_max_hcl", "Highest HCL seen per originating node.", "node_id") m.knowledgeVector = newGaugeVec(reg, "knox_knowledge_max_hcl", "Highest HCL seen per originating node.", "node_id")
+1 -1
View File
@@ -79,7 +79,7 @@ func TestMetricsScrape(t *testing.T) {
out := string(body) out := string(body)
for _, want := range []string{ for _, want := range []string{
`knox_node_info{name="testnode"`, `knox_node_info{name="testnode"`,
`knox_observations_total{source_id="git"} 2`, `knox_observations{source_id="git"} 2`,
`knox_gossip_pulls_total 0`, `knox_gossip_pulls_total 0`,
"go_goroutines", "go_goroutines",
"process_cpu_seconds_total", "process_cpu_seconds_total",
+2 -1
View File
@@ -42,7 +42,8 @@ func New(kdb *db.KnoxDB) *Server {
func (s *Server) Serve(addr string) error { func (s *Server) Serve(addr string) error {
log.Printf("[knox] web UI at http://%s", addr) log.Printf("[knox] web UI at http://%s", addr)
return http.ListenAndServe(addr, s.mux) srv := &http.Server{Addr: addr, Handler: s.mux, ReadHeaderTimeout: 10 * time.Second, IdleTimeout: 60 * time.Second}
return srv.ListenAndServe()
} }
// ─── Dashboard ─────────────────────────────────────────────── // ─── Dashboard ───────────────────────────────────────────────
+38 -16
View File
@@ -33,11 +33,11 @@ type Node struct {
// pingResponse is the anti-entropy summary returned by /v1/ping. // pingResponse is the anti-entropy summary returned by /v1/ping.
type pingResponse struct { type pingResponse struct {
NodeID string `json:"node_id"` NodeID string `json:"node_id"`
Name string `json:"name"` Name string `json:"name"`
Vector map[string]int64 `json:"vector"` // node_id → max hcl Vector map[string]int64 `json:"vector"` // node_id → max hcl
Peers []db.PeerInfo `json:"peers"` // swarm membership this node knows Peers []db.PeerInfo `json:"peers"` // swarm membership this node knows
MaxHCL *int64 `json:"max_hcl,omitempty"` MaxHCL *int64 `json:"max_hcl,omitempty"`
} }
func (n *Node) Handler() http.Handler { func (n *Node) Handler() http.Handler {
@@ -98,10 +98,16 @@ func nextCursor(rows []db.GossipObservation) int64 {
return rows[len(rows)-1].HCL return rows[len(rows)-1].HCL
} }
// maxBatchRows bounds the number of observations a peer may push in one POST.
// Pull already pages at 500 rows, so any larger batch is at best redundant and
// at worst a flood; capping keeps memory and insert work bounded.
const maxBatchRows = 1000
func (n *Node) handleBatch(w http.ResponseWriter, r *http.Request) { func (n *Node) handleBatch(w http.ResponseWriter, r *http.Request) {
r.Body = http.MaxBytesReader(w, r.Body, 4<<20) // 4 MiB
body, err := io.ReadAll(r.Body) body, err := io.ReadAll(r.Body)
if err != nil { if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest) http.Error(w, "batch too large or unreadable", http.StatusBadRequest)
return return
} }
var rows []db.GossipObservation var rows []db.GossipObservation
@@ -109,6 +115,10 @@ func (n *Node) handleBatch(w http.ResponseWriter, r *http.Request) {
http.Error(w, "bad json: "+err.Error(), http.StatusBadRequest) http.Error(w, "bad json: "+err.Error(), http.StatusBadRequest)
return return
} }
if len(rows) > maxBatchRows {
http.Error(w, fmt.Sprintf("batch too large: %d rows (max %d)", len(rows), maxBatchRows), http.StatusBadRequest)
return
}
inserted, err := n.Kdb.PushObservations(rows) inserted, err := n.Kdb.PushObservations(rows)
if err != nil { if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError) http.Error(w, err.Error(), http.StatusInternalServerError)
@@ -123,11 +133,11 @@ func (n *Node) handleBatch(w http.ResponseWriter, r *http.Request) {
// DiffSummary is the divergence snapshot served at /v1/diff: everything this // DiffSummary is the divergence snapshot served at /v1/diff: everything this
// node observes, plus the status of auto-threaded threads (for tombstones). // node observes, plus the status of auto-threaded threads (for tombstones).
type DiffSummary struct { type DiffSummary struct {
NodeID string `json:"node_id"` NodeID string `json:"node_id"`
Name string `json:"name"` Name string `json:"name"`
Fingerprints []string `json:"fingerprints"` Fingerprints []string `json:"fingerprints"`
ThreadStatus map[string]string `json:"thread_status"` ThreadStatus map[string]string `json:"thread_status"`
ObsCount int `json:"obs_count"` ObsCount int `json:"obs_count"`
} }
func (n *Node) handleDiff(w http.ResponseWriter, r *http.Request) { func (n *Node) handleDiff(w http.ResponseWriter, r *http.Request) {
@@ -163,9 +173,9 @@ func writeJSON(w http.ResponseWriter, v any) {
// Client is the pull/push half used by the exchange loop. // Client is the pull/push half used by the exchange loop.
type Client struct { type Client struct {
NodeID string NodeID string
Addr string Addr string
Timeout time.Duration Timeout time.Duration
} }
func (c *Client) Ping() (*pingResponse, error) { func (c *Client) Ping() (*pingResponse, error) {
@@ -232,6 +242,9 @@ func (c *Client) Push(kdb *db.KnoxDB, peerVector map[string]int64) (int, error)
return 0, err return 0, err
} }
defer resp.Body.Close() defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return 0, fmt.Errorf("push to %s: %s", c.Addr, resp.Status)
}
var out struct { var out struct {
Accepted int `json:"accepted"` Accepted int `json:"accepted"`
} }
@@ -273,7 +286,7 @@ func (c *Client) Diff() (*DiffSummary, error) {
// node pulls/pushes directly with every other node it learns about. // node pulls/pushes directly with every other node it learns about.
// //
// m, when non-nil, receives gossip event counters (nil for one-shot CLI runs). // m, when non-nil, receives gossip event counters (nil for one-shot CLI runs).
func Run(kdb *db.KnoxDB, m *metrics.Metrics, static []string) { func Run(kdb *db.KnoxDB, m *metrics.Metrics, static []string) int {
myID := kdb.NodeID() myID := kdb.NodeID()
// Seed the work queue with static config plus persisted discoveries. // Seed the work queue with static config plus persisted discoveries.
@@ -284,6 +297,7 @@ func Run(kdb *db.KnoxDB, m *metrics.Metrics, static []string) {
seen := make(map[string]bool) // addr → handled (also suppresses self) seen := make(map[string]bool) // addr → handled (also suppresses self)
queue := 0 queue := 0
pulledTotal := 0
for queue < len(work) { for queue < len(work) {
addr := strings.TrimSpace(work[queue]) addr := strings.TrimSpace(work[queue])
queue++ queue++
@@ -338,11 +352,18 @@ func Run(kdb *db.KnoxDB, m *metrics.Metrics, static []string) {
pulled += n pulled += n
} }
} }
pulledTotal += pulled
if m != nil { if m != nil {
m.IncrementPull(pulled) m.IncrementPull(pulled)
} }
pushed, _ := c.Push(kdb, p.Vector) pushed, err := c.Push(kdb, p.Vector)
if err != nil {
log.Printf("[gossip] push %s: %v", addr, err)
if m != nil {
m.IncrementErrors()
}
}
if m != nil { if m != nil {
m.IncrementPush(pushed) m.IncrementPush(pushed)
} }
@@ -356,6 +377,7 @@ func Run(kdb *db.KnoxDB, m *metrics.Metrics, static []string) {
log.Printf("[gossip] synced with %s (%s): in sync (pushed %d)", p.NodeID, addr, pushed) log.Printf("[gossip] synced with %s (%s): in sync (pushed %d)", p.NodeID, addr, pushed)
} }
} }
return pulledTotal
} }
func kdbVectorGet(kdb *db.KnoxDB, nodeID string) int64 { func kdbVectorGet(kdb *db.KnoxDB, nodeID string) int64 {
+126 -4
View File
@@ -1,9 +1,14 @@
package watch package watch
import ( import (
"bytes"
"encoding/json"
"net/http"
"net/http/httptest" "net/http/httptest"
"path/filepath" "path/filepath"
"strings"
"testing" "testing"
"time"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
) )
@@ -228,16 +233,133 @@ func TestGossipIdempotent(t *testing.T) {
defer sb.Close() defer sb.Close()
Run(b, nil, []string{sa.URL}) Run(b, nil, []string{sa.URL})
before, err := b.EntryCount() entryBefore, err := b.EntryCount()
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
stats, err := b.Stats()
if err != nil {
t.Fatal(err)
}
obsBefore, _ := stats["total_observations"].(int)
Run(b, nil, []string{sa.URL}) Run(b, nil, []string{sa.URL})
after, err := b.EntryCount()
entryAfter, err := b.EntryCount()
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if before != after { if entryBefore != entryAfter {
t.Errorf("second sweep changed entry count: %d -> %d", before, after) t.Errorf("second sweep changed entry count: %d -> %d", entryBefore, entryAfter)
}
stats, err = b.Stats()
if err != nil {
t.Fatal(err)
}
obsAfter, _ := stats["total_observations"].(int)
if obsBefore != obsAfter {
t.Errorf("second sweep duplicated observations: %d -> %d", obsBefore, obsAfter)
}
}
// TestHandleBatchRejectsOversizedBatch: more than maxBatchRows in one POST
// must be refused up front, before any insert work.
func TestHandleBatchRejectsOversizedBatch(t *testing.T) {
b := tmpKnoxDB(t)
node := &Node{Kdb: b, Name: "B"}
sv := httptest.NewServer(node.Handler())
defer sv.Close()
rows := make([]db.GossipObservation, maxBatchRows+1)
for i := range rows {
rows[i] = db.GossipObservation{
NodeID: "0123456789abcdef0123456789abcdef", HCL: int64(i + 1),
Fingerprint: "fp:oversized", SourceID: "test", Title: "t", Summary: "s",
CollectedAt: "2026-08-29T00:00:00Z",
}
}
body, _ := json.Marshal(rows)
resp, err := http.Post(sv.URL+"/v1/obs/batch", "application/json", bytes.NewReader(body))
if err != nil {
t.Fatalf("post: %v", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusBadRequest {
t.Errorf("oversized batch: want 400, got %d", resp.StatusCode)
}
}
// TestHandleBatchRejectsHugeBody: an oversized body (beyond the 4 MiB cap)
// must be refused even when the row count is small.
func TestHandleBatchRejectsHugeBody(t *testing.T) {
b := tmpKnoxDB(t)
node := &Node{Kdb: b, Name: "B"}
sv := httptest.NewServer(node.Handler())
defer sv.Close()
rows := []db.GossipObservation{{
NodeID: "0123456789abcdef0123456789abcdef", HCL: 1,
Fingerprint: "fp:huge", SourceID: "test", Title: "t",
Summary: strings.Repeat("x", 5<<20), // 5 MiB summary
CollectedAt: "2026-08-29T00:00:00Z",
}}
body, _ := json.Marshal(rows)
resp, err := http.Post(sv.URL+"/v1/obs/batch", "application/json", bytes.NewReader(body))
if err != nil {
t.Fatalf("post: %v", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusBadRequest {
t.Errorf("huge body: want 400, got %d", resp.StatusCode)
}
}
// TestHandleBatchSkipsSelfRows: rows claiming the receiver's own node_id are
// dropped at the HTTP layer too (the poisoning vector), reported as conflicts.
func TestHandleBatchSkipsSelfRows(t *testing.T) {
b := tmpKnoxDB(t)
node := &Node{Kdb: b, Name: "B"}
sv := httptest.NewServer(node.Handler())
defer sv.Close()
rows := []db.GossipObservation{
{NodeID: b.NodeID(), HCL: 1, Fingerprint: "fp:self1", SourceID: "test", Title: "t", Summary: "s", CollectedAt: "2026-08-29T00:00:00Z"},
{NodeID: b.NodeID(), HCL: 2, Fingerprint: "fp:self2", SourceID: "test", Title: "t", Summary: "s", CollectedAt: "2026-08-29T00:00:00Z"},
}
body, _ := json.Marshal(rows)
resp, err := http.Post(sv.URL+"/v1/obs/batch", "application/json", bytes.NewReader(body))
if err != nil {
t.Fatalf("post: %v", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("want 200, got %d", resp.StatusCode)
}
var out struct {
Accepted int `json:"accepted"`
Conflict int `json:"conflict"`
}
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
t.Fatalf("decode: %v", err)
}
if out.Accepted != 0 || out.Conflict != 2 {
t.Errorf("self rows: want accepted=0 conflict=2, got accepted=%d conflict=%d", out.Accepted, out.Conflict)
}
}
// TestClientPushSurfacesHTTPError: a non-2xx push response must surface as an
// error — silently treating it as accepted=0 would hide a broken sync direction.
func TestClientPushSurfacesHTTPError(t *testing.T) {
a := tmpKnoxDB(t)
seedObs(a, "AAA")
sv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, "boom", http.StatusInternalServerError)
}))
defer sv.Close()
c := &Client{Addr: sv.URL, Timeout: 5 * time.Second}
if _, err := c.Push(a, map[string]int64{}); err == nil {
t.Fatal("expected push error on HTTP 500, got nil")
} }
} }
+141 -67
View File
@@ -1,11 +1,14 @@
package watch package watch
import ( import (
"io/fs"
"log" "log"
"net/http" "net/http"
"os" "os"
"path/filepath" "path/filepath"
"strings" "strings"
"sync"
"sync/atomic"
"time" "time"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
@@ -28,6 +31,44 @@ type Watcher struct {
metrics *metrics.Metrics metrics *metrics.Metrics
} }
// fileDebouncer schedules one ingest per path a settle-window after the last
// relevant event (trailing edge): writers that emit bursts (multi-write appends,
// atomic save = temp-write + rename) settle before anything is read, so a
// partial file is never recorded as final state. Timers are removed once they
// fire, keeping the map bounded by recently-active paths.
type fileDebouncer struct {
mu sync.Mutex
timers map[string]*time.Timer
}
func newFileDebouncer() *fileDebouncer {
return &fileDebouncer{timers: make(map[string]*time.Timer)}
}
func (d *fileDebouncer) schedule(path string, delay time.Duration, fn func(string)) {
d.mu.Lock()
defer d.mu.Unlock()
if t, ok := d.timers[path]; ok {
t.Stop()
}
d.timers[path] = time.AfterFunc(delay, func() {
d.mu.Lock()
delete(d.timers, path)
d.mu.Unlock()
fn(path)
})
}
// cancel drops any pending ingest for path (e.g. the file was deleted).
func (d *fileDebouncer) cancel(path string) {
d.mu.Lock()
defer d.mu.Unlock()
if t, ok := d.timers[path]; ok {
t.Stop()
delete(d.timers, path)
}
}
func New(kdb *db.KnoxDB, dirs []string) *Watcher { func New(kdb *db.KnoxDB, dirs []string) *Watcher {
// Detect Obsidian vault // Detect Obsidian vault
vault, _ := ingest.DetectObsidianVault() vault, _ := ingest.DetectObsidianVault()
@@ -57,7 +98,7 @@ func (w *Watcher) Start() error {
// seed is still ingesting. Vault/dirs are logged after the seed below. // seed is still ingesting. Vault/dirs are logged after the seed below.
gossipAddr := ListenAddr() gossipAddr := ListenAddr()
node := &Node{Kdb: w.knoxDB, Name: "knox", Addr: gossipAddr, Metrics: w.metrics} node := &Node{Kdb: w.knoxDB, Name: "knox", Addr: gossipAddr, Metrics: w.metrics}
srv := &http.Server{Addr: gossipAddr, Handler: node.Handler()} srv := &http.Server{Addr: gossipAddr, Handler: node.Handler(), ReadHeaderTimeout: 10 * time.Second, IdleTimeout: 60 * time.Second}
go func() { go func() {
log.Printf("[knox] gossip listening on %s", gossipAddr) log.Printf("[knox] gossip listening on %s", gossipAddr)
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
@@ -67,7 +108,7 @@ func (w *Watcher) Start() error {
// Prometheus scraping on a dedicated port (KNOX_METRICS_ADDR). // Prometheus scraping on a dedicated port (KNOX_METRICS_ADDR).
metricsAddr := MetricsAddr() metricsAddr := MetricsAddr()
metricsSrv := &http.Server{Addr: metricsAddr, Handler: node.Metrics.Handler()} metricsSrv := &http.Server{Addr: metricsAddr, Handler: node.Metrics.Handler(), ReadHeaderTimeout: 10 * time.Second, IdleTimeout: 60 * time.Second}
go func() { go func() {
log.Printf("[knox] metrics listening on %s", metricsAddr) log.Printf("[knox] metrics listening on %s", metricsAddr)
if err := metricsSrv.ListenAndServe(); err != nil && err != http.ErrServerClosed { if err := metricsSrv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
@@ -86,17 +127,13 @@ func (w *Watcher) Start() error {
log.Printf("[knox] obsidian vault: %s", w.vault) log.Printf("[knox] obsidian vault: %s", w.vault)
} }
debounceMap := make(map[string]time.Time) debounce := newFileDebouncer()
browserTicker := time.NewTicker(browserInterval) browserTicker := time.NewTicker(browserInterval)
giteaTicker := time.NewTicker(10 * time.Minute) giteaTicker := time.NewTicker(10 * time.Minute)
gitTicker := time.NewTicker(10 * time.Minute) gitTicker := time.NewTicker(10 * time.Minute)
threadTicker := time.NewTicker(10 * time.Minute) threadTicker := time.NewTicker(10 * time.Minute)
gossipTicker := time.NewTicker(gossipInterval) gossipTicker := time.NewTicker(gossipInterval)
browserRunning := false var browserRunning, giteaRunning, gitRunning, threadRunning, gossipRunning atomic.Bool
giteaRunning := false
gitRunning := false
threadRunning := false
gossipRunning := false
for { for {
select { select {
@@ -104,63 +141,76 @@ func (w *Watcher) Start() error {
if !ok { if !ok {
return nil return nil
} }
if !w.isRelevantEvent(event) {
continue // fsnotify is non-recursive: files appearing inside newly created
// subdirectories would otherwise be invisible to the daemon.
isDir := false
if event.Has(fsnotify.Create) {
if info, err := os.Stat(event.Name); err == nil && info.IsDir() {
isDir = true
if err := watcher.Add(event.Name); err != nil {
log.Printf("[knox] cannot watch new dir %s: %v", event.Name, err)
} else {
log.Printf("[knox] watching new dir %s", event.Name)
}
}
} }
now := time.Now() if event.Has(fsnotify.Remove) {
if last, ok := debounceMap[event.Name]; ok && now.Sub(last) < w.debounce { // Drop any pending ingest for a deleted file. (No tombstone is
// written yet — the entry lingers until reconcile/prune.)
debounce.cancel(event.Name)
continue
}
if isDir || !w.isRelevantEvent(event) {
continue continue
} }
debounceMap[event.Name] = now
trigger := eventOpName(event.Op) trigger := eventOpName(event.Op)
log.Printf("[knox] %s %s", trigger, filepath.Base(event.Name)) log.Printf("[knox] %s %s", trigger, filepath.Base(event.Name))
w.ingestFile(event.Name, trigger) name := event.Name
debounce.schedule(name, w.debounce, func(path string) {
w.ingestFile(path, trigger)
})
case <-browserTicker.C: case <-browserTicker.C:
if browserRunning { if !browserRunning.CompareAndSwap(false, true) {
continue continue
} }
browserRunning = true
go func() { go func() {
defer func() { browserRunning = false }() defer browserRunning.Store(false)
w.ingestBrowserHistory() w.ingestBrowserHistory()
}() }()
case <-giteaTicker.C: case <-giteaTicker.C:
if giteaRunning { if !giteaRunning.CompareAndSwap(false, true) {
continue continue
} }
giteaRunning = true
go func() { go func() {
defer func() { giteaRunning = false }() defer giteaRunning.Store(false)
w.ingestGitea() w.ingestGitea()
}() }()
case <-gitTicker.C: case <-gitTicker.C:
if gitRunning { if !gitRunning.CompareAndSwap(false, true) {
continue continue
} }
gitRunning = true
go func() { go func() {
defer func() { gitRunning = false }() defer gitRunning.Store(false)
w.ingestGit() w.ingestGit()
}() }()
case <-threadTicker.C: case <-threadTicker.C:
if threadRunning { if !threadRunning.CompareAndSwap(false, true) {
continue continue
} }
threadRunning = true
go func() { go func() {
defer func() { threadRunning = false }() defer threadRunning.Store(false)
w.autoThread() w.autoThread()
}() }()
case <-gossipTicker.C: case <-gossipTicker.C:
if gossipRunning { if !gossipRunning.CompareAndSwap(false, true) {
continue continue
} }
gossipRunning = true
go func() { go func() {
defer func() { gossipRunning = false }() defer gossipRunning.Store(false)
w.syncGossip() w.syncGossip()
}() }()
@@ -175,44 +225,38 @@ func (w *Watcher) Start() error {
func (w *Watcher) seed(watcher *fsnotify.Watcher) { func (w *Watcher) seed(watcher *fsnotify.Watcher) {
for _, dir := range w.dirs { for _, dir := range w.dirs {
abs, _ := filepath.Abs(dir) if err := w.watchTree(watcher, dir); err != nil {
if err := watcher.Add(abs); err != nil { log.Printf("[knox] cannot watch %s: %v", dir, err)
log.Printf("[knox] cannot watch %s: %v", abs, err)
continue
} }
log.Printf("[knox] watching %s", abs)
} }
// Watch Obsidian vault // Watch Obsidian vault (every subdirectory, non-recursively mirrored)
if w.vault != "" { if w.vault != "" {
if err := watcher.Add(w.vault); err != nil { if err := w.watchTree(watcher, w.vault); err != nil {
log.Printf("[knox] cannot watch obsidian vault %s: %v", w.vault, err) log.Printf("[knox] cannot watch obsidian vault %s: %v", w.vault, err)
} else {
log.Printf("[knox] watching %s (obsidian)", w.vault)
} }
} }
// Seed existing files // Seed existing files anywhere under the watched dirs. fsnotify watches the
for _, ing := range w.fileIngesters { // whole tree, so live events cover any depth; this walk covers startup so
for _, dir := range w.dirs { // pre-existing nested files (e.g. skills at two+ levels) are indexed too.
patterns := []string{ for _, dir := range w.dirs {
filepath.Join(dir, "*"), _ = filepath.WalkDir(dir, func(path string, d fs.DirEntry, err error) error {
filepath.Join(dir, "*", "SKILL.md"), if err != nil || d.IsDir() {
return nil // skip unreadable entries; dirs are handled by the watch
} }
for _, pattern := range patterns { for _, ing := range w.fileIngesters {
entries, _ := filepath.Glob(pattern) if MatchesIngester(path, ing.SourceID()) {
for _, path := range entries {
if !MatchesIngester(path, ing.SourceID()) {
continue
}
w.ingestFileWith(path, ing, "seed") w.ingestFileWith(path, ing, "seed")
return nil
} }
} }
} return nil
})
} }
// Seed Obsidian notes via the full walk: skips dot-dirs (.trash, .obsidian) // Seed Obsidian notes via the full walk: skips dot-dirs (.trash, .obsidian)
// and covers all depths — glob patterns would match dot-dirs and miss depth >2. // and covers all depths.
if w.vault != "" { if w.vault != "" {
if notes, err := ingest.NewObsidianIngester(w.vault).IngestAll(); err == nil { if notes, err := ingest.NewObsidianIngester(w.vault).IngestAll(); err == nil {
for _, r := range notes { for _, r := range notes {
@@ -222,6 +266,39 @@ func (w *Watcher) seed(watcher *fsnotify.Watcher) {
} }
} }
// watchTree adds a directory and every non-hidden subdirectory to the watcher,
// mirroring fsnotify's non-recursive API with an explicit walk. Hidden
// directories (.git, .obsidian, .trash) are skipped so their churn doesn't burn
// inotify watches.
func (w *Watcher) watchTree(watcher *fsnotify.Watcher, root string) error {
rootAbs, err := filepath.Abs(root)
if err != nil {
return err
}
if err := watcher.Add(rootAbs); err != nil {
return err
}
log.Printf("[knox] watching %s", rootAbs)
return filepath.WalkDir(rootAbs, func(path string, d fs.DirEntry, err error) error {
if err != nil {
return nil
}
if !d.IsDir() {
return nil
}
if path != rootAbs && strings.HasPrefix(d.Name(), ".") {
return filepath.SkipDir
}
if path == rootAbs {
return nil
}
if err := watcher.Add(path); err != nil {
log.Printf("[knox] cannot watch %s: %v", path, err)
}
return nil
})
}
func (w *Watcher) ingestFile(path, trigger string) { func (w *Watcher) ingestFile(path, trigger string) {
for _, ing := range w.fileIngesters { for _, ing := range w.fileIngesters {
if MatchesIngester(path, ing.SourceID()) { if MatchesIngester(path, ing.SourceID()) {
@@ -290,8 +367,9 @@ func (w *Watcher) recordResult(result *ingest.IngestResult, trigger string) {
if result.SourceID == "opencode-session" { if result.SourceID == "opencode-session" {
sessionID, _ := result.Provenance["session_id"].(string) sessionID, _ := result.Provenance["session_id"].(string)
if sessionID != "" { if sessionID != "" {
status := "active" if err := w.knoxDB.UpsertSession(sessionID, result.Project, result.Title, "active"); err != nil {
w.knoxDB.UpsertSession(sessionID, result.Project, result.Title, status) log.Printf("[knox] session upsert error: %v", err)
}
} }
} }
} }
@@ -456,29 +534,23 @@ func (w *Watcher) syncGossip() {
return return
} }
before := 0 pulled := Run(w.knoxDB, w.metrics, peers)
if n, err := w.knoxDB.EntryCount(); err == nil {
before = n
}
Run(w.knoxDB, w.metrics, peers) // A pull only appends to the observation log; the entries cache and
// auto-threads are derived state that reconcile rebuilds. Comparing entry
// If new observations arrived, reconcile to pick up entries/threads they // counts can never trigger this (pushes/pulls never touch entries directly),
// imply (deterministic log → derived rebuild). // so reconcile fires on the sweep's newly-inserted observation count.
if after, err := w.knoxDB.EntryCount(); err == nil && after > before { if pulled > 0 {
created, linked, err := Reconcile(w.knoxDB) created, linked, err := Reconcile(w.knoxDB)
if err != nil { if err != nil {
log.Printf("[knox] gossip reconcile: %v", err) log.Printf("[knox] gossip reconcile: %v", err)
return return
} }
log.Printf("[knox] gossip reconcile done: %d created, %d linked", created, linked) log.Printf("[knox] gossip reconcile done: %d created, %d linked after pulling %d obs", created, linked, pulled)
} }
} }
func (w *Watcher) isRelevantEvent(event fsnotify.Event) bool { func (w *Watcher) isRelevantEvent(event fsnotify.Event) bool {
if event.Has(fsnotify.Remove) || event.Has(fsnotify.Rename) {
return false
}
name := filepath.Base(event.Name) name := filepath.Base(event.Name)
// Opencode session diffs and logs // Opencode session diffs and logs
if strings.HasPrefix(name, "ses_") || strings.HasSuffix(name, ".log") { if strings.HasPrefix(name, "ses_") || strings.HasSuffix(name, ".log") {
@@ -512,6 +584,8 @@ func eventOpName(op fsnotify.Op) string {
return "inotify:WRITE" return "inotify:WRITE"
case op.Has(fsnotify.Chmod): case op.Has(fsnotify.Chmod):
return "inotify:CHMOD" return "inotify:CHMOD"
case op.Has(fsnotify.Rename):
return "inotify:RENAME"
default: default:
return "inotify:UNKNOWN" return "inotify:UNKNOWN"
} }
+9 -2
View File
@@ -1,15 +1,17 @@
package main package main
import ( import (
"context"
"errors"
"fmt" "fmt"
"log" "log"
"os" "os"
mcpServer "github.com/mark3labs/mcp-go/server"
"github.com/spf13/cobra"
knoxcmd "github.com/david/knox/internal/cmd" knoxcmd "github.com/david/knox/internal/cmd"
"github.com/david/knox/internal/db" "github.com/david/knox/internal/db"
knoxserver "github.com/david/knox/internal/server" knoxserver "github.com/david/knox/internal/server"
mcpServer "github.com/mark3labs/mcp-go/server"
"github.com/spf13/cobra"
) )
func newServeCmd(kdb *db.KnoxDB) *cobra.Command { func newServeCmd(kdb *db.KnoxDB) *cobra.Command {
@@ -75,6 +77,11 @@ and maintains a searchable index. Use 'knox watch' for daemon mode.`,
s := knoxcmd.NewMCPServer(kdb) s := knoxcmd.NewMCPServer(kdb)
log.Printf("[knox] starting MCP stdio server") log.Printf("[knox] starting MCP stdio server")
if err := mcpServer.ServeStdio(s); err != nil { if err := mcpServer.ServeStdio(s); err != nil {
// SIGINT/SIGTERM cancels the server context: clean shutdown,
// not a failure.
if errors.Is(err, context.Canceled) {
return nil
}
return err return err
} }
return nil return nil