feat: write leafs incrementally without rebuilding the graph. (#29)
Ladybug 0.19 stays FTS/HNSW queryable on MERGE of new ids; DROP INDEX was the fatal path. bin/brain/add.go and POST /ingest land facts+info in one transaction so watch/mail/git can become leafs now (Gitea #14).
This commit is contained in:
@@ -144,7 +144,15 @@ func (s *Server) mcpCall(r *http.Request, params json.RawMessage) (any, error) {
|
||||
return nil, fmt.Errorf("cancelled")
|
||||
}
|
||||
defer s.release()
|
||||
body, err = s.api.Ingest(r.Context())
|
||||
var payload []byte
|
||||
text := strings.TrimSpace(fmt.Sprint(p.Arguments["text"]))
|
||||
if text != "" && text != "<nil>" {
|
||||
payload, err = json.Marshal(p.Arguments)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
body, err = s.api.Ingest(r.Context(), payload)
|
||||
default:
|
||||
return nil, fmt.Errorf("unknown tool %s", p.Name)
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
@@ -27,7 +28,7 @@ type API interface {
|
||||
Get(ctx context.Context, id string, body bool) ([]byte, error)
|
||||
Stats(ctx context.Context) ([]byte, error)
|
||||
Audit(ctx context.Context) ([]byte, error)
|
||||
Ingest(ctx context.Context) ([]byte, error)
|
||||
Ingest(ctx context.Context, body []byte) ([]byte, error)
|
||||
}
|
||||
|
||||
type Server struct {
|
||||
@@ -59,7 +60,7 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
case PathAudit:
|
||||
s.handleJSON(w, r, s.api.Audit)
|
||||
case PathIngest:
|
||||
s.handleJSON(w, r, s.api.Ingest)
|
||||
s.handleIngest(w, r)
|
||||
case PathOpenAPI:
|
||||
s.handleOpenAPI(w, r)
|
||||
case PathMCP:
|
||||
@@ -116,6 +117,24 @@ func (s *Server) handleJSON(w http.ResponseWriter, r *http.Request, fn func(cont
|
||||
writeAPI(w, body, err)
|
||||
}
|
||||
|
||||
func (s *Server) handleIngest(w http.ResponseWriter, r *http.Request) {
|
||||
var raw []byte
|
||||
if r.Method == http.MethodPost {
|
||||
b, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
|
||||
if err != nil {
|
||||
writeJSON(w, http.StatusBadRequest, map[string]any{"error": "read body"})
|
||||
return
|
||||
}
|
||||
raw = b
|
||||
}
|
||||
if !s.acquire(w, r) {
|
||||
return
|
||||
}
|
||||
defer s.release()
|
||||
body, err := s.api.Ingest(r.Context(), raw)
|
||||
writeAPI(w, body, err)
|
||||
}
|
||||
|
||||
func (s *Server) tryAcquire(r *http.Request) bool {
|
||||
return s.acquire(nopWriter{}, r)
|
||||
}
|
||||
@@ -191,11 +210,30 @@ func (ExecSearcher) Get(context.Context, string, bool) ([]byte, error) {
|
||||
}
|
||||
func (ExecSearcher) Stats(context.Context) ([]byte, error) { return nil, errUnimplemented }
|
||||
func (ExecSearcher) Audit(context.Context) ([]byte, error) { return nil, errUnimplemented }
|
||||
func (ExecSearcher) Ingest(context.Context) ([]byte, error) {
|
||||
return json.Marshal(map[string]any{
|
||||
"mode": "rebuild",
|
||||
"command": "bin/brain/index.go --rebuild",
|
||||
})
|
||||
func (b ExecSearcher) Ingest(ctx context.Context, body []byte) ([]byte, error) {
|
||||
if len(strings.TrimSpace(string(body))) == 0 {
|
||||
return json.Marshal(map[string]any{
|
||||
"mode": "add",
|
||||
"command": "bin/brain/add.go",
|
||||
"rebuild": "bin/brain/index.go --rebuild",
|
||||
})
|
||||
}
|
||||
root := os.Getenv("KB_ROOT")
|
||||
if root == "" {
|
||||
root = "."
|
||||
}
|
||||
cmd := exec.CommandContext(ctx, filepath.Join(root, "bin/kb/add"), "--json")
|
||||
cmd.Stdin = strings.NewReader(string(body))
|
||||
cmd.Dir = root
|
||||
out, err := cmd.Output()
|
||||
if err != nil {
|
||||
var exitErr *exec.ExitError
|
||||
if errors.As(err, &exitErr) {
|
||||
return nil, errors.New("add failed: " + strings.TrimSpace(string(exitErr.Stderr)))
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func defaultSearchCmd(root string) string {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package httpapi
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
@@ -64,8 +65,11 @@ func (f *fakeSearcher) Audit(context.Context) ([]byte, error) {
|
||||
return []byte(`{"status":"ok"}`), nil
|
||||
}
|
||||
|
||||
func (f *fakeSearcher) Ingest(context.Context) ([]byte, error) {
|
||||
return []byte(`{"mode":"rebuild","command":"bin/brain/index.go --rebuild"}`), nil
|
||||
func (f *fakeSearcher) Ingest(_ context.Context, body []byte) ([]byte, error) {
|
||||
if len(bytes.TrimSpace(body)) == 0 {
|
||||
return []byte(`{"mode":"add","command":"bin/brain/add.go"}`), nil
|
||||
}
|
||||
return []byte(`{"mode":"add","ids":["fake-leaf"]}`), nil
|
||||
}
|
||||
|
||||
func (f *fakeSearcher) count() int {
|
||||
@@ -193,6 +197,27 @@ func TestStatsAuditIngest(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngestIsAddNotRebuildHint(t *testing.T) {
|
||||
h := NewServer(&fakeSearcher{}, 1)
|
||||
code, body := get(t, h, "/ingest")
|
||||
if code != http.StatusOK {
|
||||
t.Fatalf("GET /ingest code = %d body=%s", code, body)
|
||||
}
|
||||
if strings.Contains(string(body), `"add":"v2"`) || strings.Contains(string(body), "write is v2") {
|
||||
t.Fatalf("GET /ingest still a v2 hint: %s", body)
|
||||
}
|
||||
if !strings.Contains(string(body), "bin/brain/add.go") {
|
||||
t.Fatalf("GET /ingest should name add.go: %s", body)
|
||||
}
|
||||
code, body = postJSON(t, h, "/ingest", `{"text":"hello","root":"info","source":"t"}`)
|
||||
if code != http.StatusOK {
|
||||
t.Fatalf("POST /ingest code = %d body=%s", code, body)
|
||||
}
|
||||
if !strings.Contains(string(body), "fake-leaf") {
|
||||
t.Fatalf("POST /ingest should add: %s", body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHTTPPackageDoesNotExecPython(t *testing.T) {
|
||||
raw, err := os.ReadFile("server.go")
|
||||
if err != nil {
|
||||
|
||||
@@ -47,7 +47,15 @@ var Ops = []Op{
|
||||
},
|
||||
{Path: PathStats, Method: "get", ID: "stats", Summary: "index health", MCP: true},
|
||||
{Path: PathAudit, Method: "get", ID: "audit", Summary: "facts confidence histogram", MCP: true},
|
||||
{Path: PathIngest, Method: "get", ID: "ingest", Summary: "rebuild hint (write is v2)", MCP: true},
|
||||
{
|
||||
Path: PathIngest, Method: "post", ID: "ingest", Summary: "add a leaf without rebuild",
|
||||
MCP: true,
|
||||
Params: []Param{
|
||||
{Name: "text", In: "query", Type: "string", Description: "leaf text (omit for CLI hint)"},
|
||||
{Name: "root", In: "query", Type: "string", Description: "facts or info (default info)"},
|
||||
{Name: "source", In: "query", Type: "string", Description: "evidence pointer; facts need two sources"},
|
||||
},
|
||||
},
|
||||
{Path: PathOpenAPI, Method: "get", ID: "openapi", Summary: "OpenAPI 3 document for this server"},
|
||||
}
|
||||
|
||||
|
||||
@@ -38,11 +38,24 @@ func TestMCPToolsMatchOpenAPIPaths(t *testing.T) {
|
||||
t.Fatalf("MCP tool %s has no OpenAPI path %s", tool.Name, path)
|
||||
}
|
||||
}
|
||||
for _, need := range []string{"search", "get", "stats", "audit"} {
|
||||
for _, need := range []string{"search", "get", "stats", "audit", "ingest"} {
|
||||
if !names[need] {
|
||||
t.Fatalf("MCP tools missing %s: %v", need, names)
|
||||
}
|
||||
}
|
||||
var ingest MCPTool
|
||||
for _, tool := range tools {
|
||||
if tool.Name == "ingest" {
|
||||
ingest = tool
|
||||
break
|
||||
}
|
||||
}
|
||||
if strings.Contains(ingest.Description, "v2") {
|
||||
t.Fatalf("ingest still a v2 hint: %s", ingest.Description)
|
||||
}
|
||||
if !strings.Contains(ingest.Description, "add") {
|
||||
t.Fatalf("ingest should describe add: %s", ingest.Description)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenAPIHTTP(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user