Skip to content

Commit 40c796e

Browse files
committed
feat(rag): add rag-index-status worker and webhook trigger
1 parent ae7d871 commit 40c796e

3 files changed

Lines changed: 307 additions & 0 deletions

File tree

model/rag/webhook.go

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,46 @@
1+
package rag
2+
3+
import (
4+
"errors"
5+
6+
"github.com/cozy/cozy-stack/model/instance"
7+
"github.com/cozy/cozy-stack/model/job"
8+
"github.com/cozy/cozy-stack/pkg/couchdb"
9+
)
10+
11+
// ragStatusTriggerID is the fixed ID of the per-instance @webhook trigger that
12+
// feeds the "rag-index-status" worker, so it can be looked up in a single query.
13+
const ragStatusTriggerID = "rag-index-status"
14+
15+
// EnsureRAGWebhook returns the URL of the per-instance @webhook trigger feeding
16+
// the "rag-index-status" worker, creating it on the first call.
17+
func EnsureRAGWebhook(inst *instance.Instance) (string, error) {
18+
sched := job.System()
19+
20+
t, err := sched.GetTrigger(inst, ragStatusTriggerID)
21+
if err != nil {
22+
if !errors.Is(err, job.ErrNotFoundTrigger) {
23+
return "", err
24+
}
25+
t, err = job.NewTrigger(inst, job.TriggerInfos{
26+
TID: ragStatusTriggerID,
27+
Type: "@webhook",
28+
WorkerType: "rag-index-status",
29+
}, nil)
30+
if err != nil {
31+
return "", err
32+
}
33+
if err = sched.AddTrigger(t); err != nil {
34+
if !couchdb.IsConflictError(err) {
35+
return "", err
36+
}
37+
// A concurrent first call already created the trigger; reuse it.
38+
if t, err = sched.GetTrigger(inst, ragStatusTriggerID); err != nil {
39+
return "", err
40+
}
41+
} else {
42+
inst.Logger().WithNamespace("rag").Infof("RAG webhook trigger created: %s", t.ID())
43+
}
44+
}
45+
return inst.PageURL("/jobs/webhooks/"+t.ID(), nil), nil
46+
}

worker/rag/callback_status.go

Lines changed: 84 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,84 @@
1+
package rag
2+
3+
import (
4+
"encoding/json"
5+
"fmt"
6+
"runtime"
7+
"time"
8+
9+
"github.com/cozy/cozy-stack/model/job"
10+
modelrag "github.com/cozy/cozy-stack/model/rag"
11+
"github.com/cozy/cozy-stack/pkg/couchdb"
12+
)
13+
14+
func init() {
15+
job.AddWorker(&job.WorkerConfig{
16+
WorkerType: "rag-index-status",
17+
Concurrency: runtime.NumCPU(),
18+
MaxExecCount: 2,
19+
Timeout: 30 * time.Second,
20+
WorkerFunc: WorkerIndexStatus,
21+
})
22+
}
23+
24+
type statusMessage struct {
25+
Partition string `json:"partition"`
26+
FileID string `json:"file_id"`
27+
Status string `json:"status"` // "success" | "error" | "notsupported"
28+
Timestamp string `json:"timestamp"` // RFC3339Nano
29+
// Metadata is the file metadata carried by the callback. Its "version" field
30+
// holds the content md5sum, used to drop callbacks about outdated content.
31+
Metadata struct {
32+
Version string `json:"version"`
33+
} `json:"metadata"`
34+
}
35+
36+
func WorkerIndexStatus(ctx *job.TaskContext) error {
37+
inst := ctx.Instance
38+
log := inst.Logger().WithNamespace("rag")
39+
40+
raw, err := ctx.UnmarshalPayload()
41+
if err != nil {
42+
return err
43+
}
44+
data, err := json.Marshal(raw)
45+
if err != nil {
46+
return err
47+
}
48+
var msg statusMessage
49+
if err := json.Unmarshal(data, &msg); err != nil {
50+
return err
51+
}
52+
if msg.FileID == "" {
53+
return fmt.Errorf("rag-index-status: missing file_id in payload")
54+
}
55+
56+
switch msg.Status {
57+
case modelrag.RAGStatusSuccess, modelrag.RAGStatusError, modelrag.RAGStatusNotSupported:
58+
default:
59+
return fmt.Errorf("rag-index-status: unknown status %q for file %s", msg.Status, msg.FileID)
60+
}
61+
62+
var ts time.Time
63+
if msg.Timestamp != "" {
64+
ts, err = time.Parse(time.RFC3339Nano, msg.Timestamp)
65+
if err != nil {
66+
log.Warnf("rag-index-status: invalid timestamp %q for file %s, using now", msg.Timestamp, msg.FileID)
67+
ts = time.Now()
68+
}
69+
} else {
70+
log.Warnf("rag-index-status: missing timestamp for file %s, using now", msg.FileID)
71+
ts = time.Now()
72+
}
73+
74+
log.Debugf("rag-index-status: file %s status=%s ts=%s", msg.FileID, msg.Status, ts)
75+
76+
if err := modelrag.SetRAGStatus(inst, msg.FileID, msg.Status, msg.Metadata.Version, ts); err != nil {
77+
if couchdb.IsNotFoundError(err) {
78+
log.Debugf("rag-index-status: file %s not found (possibly deleted), skipping", msg.FileID)
79+
return nil
80+
}
81+
return err
82+
}
83+
return nil
84+
}

worker/rag/callback_status_test.go

Lines changed: 177 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,177 @@
1+
package rag
2+
3+
import (
4+
"encoding/json"
5+
"testing"
6+
"time"
7+
8+
"github.com/cozy/cozy-stack/model/instance/lifecycle"
9+
"github.com/cozy/cozy-stack/model/job"
10+
modelrag "github.com/cozy/cozy-stack/model/rag"
11+
"github.com/cozy/cozy-stack/model/vfs"
12+
"github.com/cozy/cozy-stack/pkg/config/config"
13+
"github.com/cozy/cozy-stack/tests/testutils"
14+
"github.com/stretchr/testify/require"
15+
)
16+
17+
func TestWorkerIndexStatus(t *testing.T) {
18+
config.UseTestFile(t)
19+
setup := testutils.NewSetup(t, "rag_status_test")
20+
inst := setup.GetTestInstance(&lifecycle.Options{})
21+
22+
runWorker := func(t *testing.T, payload map[string]interface{}) error {
23+
t.Helper()
24+
data, err := json.Marshal(payload)
25+
require.NoError(t, err)
26+
j := &job.Job{
27+
Domain: inst.Domain,
28+
WorkerType: "rag-index-status",
29+
Payload: job.Payload(data),
30+
}
31+
ctx, cancel := job.NewTaskContext("test", j, inst)
32+
defer cancel()
33+
return WorkerIndexStatus(ctx)
34+
}
35+
36+
t.Run("success status sets Status=success and Indexed=true", func(t *testing.T) {
37+
fs := inst.VFS()
38+
doc := createStatusTestFile(t, fs, "rag-success.txt")
39+
defer destroyStatusTestFile(t, fs, doc)
40+
41+
ts := time.Now().UTC().Truncate(time.Second)
42+
err := runWorker(t, map[string]interface{}{
43+
"partition": inst.Domain,
44+
"file_id": doc.DocID,
45+
"status": "success",
46+
"timestamp": ts.Format(time.RFC3339Nano),
47+
})
48+
require.NoError(t, err)
49+
50+
updated, err := fs.FileByID(doc.DocID)
51+
require.NoError(t, err)
52+
require.NotNil(t, updated.CozyMetadata.RAG)
53+
require.True(t, updated.CozyMetadata.RAG.Indexed)
54+
require.Equal(t, modelrag.RAGStatusSuccess, updated.CozyMetadata.RAG.Status)
55+
require.NotNil(t, updated.CozyMetadata.RAG.LastSuccessDate)
56+
require.Nil(t, updated.CozyMetadata.RAG.LastErrorDate)
57+
})
58+
59+
t.Run("error status sets Status=error and preserves Indexed", func(t *testing.T) {
60+
fs := inst.VFS()
61+
doc := createStatusTestFile(t, fs, "rag-error.txt")
62+
defer destroyStatusTestFile(t, fs, doc)
63+
64+
ts := time.Now().UTC().Truncate(time.Second)
65+
err := runWorker(t, map[string]interface{}{
66+
"partition": inst.Domain,
67+
"file_id": doc.DocID,
68+
"status": "error",
69+
"timestamp": ts.Format(time.RFC3339Nano),
70+
})
71+
require.NoError(t, err)
72+
73+
updated, err := fs.FileByID(doc.DocID)
74+
require.NoError(t, err)
75+
require.NotNil(t, updated.CozyMetadata.RAG)
76+
require.False(t, updated.CozyMetadata.RAG.Indexed)
77+
require.Equal(t, modelrag.RAGStatusError, updated.CozyMetadata.RAG.Status)
78+
require.Nil(t, updated.CozyMetadata.RAG.LastSuccessDate)
79+
require.NotNil(t, updated.CozyMetadata.RAG.LastErrorDate)
80+
})
81+
82+
t.Run("notsupported status sets Status=notsupported without touching Indexed or dates", func(t *testing.T) {
83+
fs := inst.VFS()
84+
doc := createStatusTestFile(t, fs, "rag-notsupported.txt")
85+
defer destroyStatusTestFile(t, fs, doc)
86+
87+
err := runWorker(t, map[string]interface{}{
88+
"partition": inst.Domain,
89+
"file_id": doc.DocID,
90+
"status": "notsupported",
91+
"timestamp": time.Now().UTC().Format(time.RFC3339Nano),
92+
})
93+
require.NoError(t, err)
94+
95+
updated, err := fs.FileByID(doc.DocID)
96+
require.NoError(t, err)
97+
require.NotNil(t, updated.CozyMetadata.RAG)
98+
require.Equal(t, modelrag.RAGStatusNotSupported, updated.CozyMetadata.RAG.Status)
99+
require.False(t, updated.CozyMetadata.RAG.Indexed)
100+
require.Nil(t, updated.CozyMetadata.RAG.LastSuccessDate)
101+
require.Nil(t, updated.CozyMetadata.RAG.LastErrorDate)
102+
})
103+
104+
t.Run("unknown status returns error", func(t *testing.T) {
105+
err := runWorker(t, map[string]interface{}{
106+
"partition": inst.Domain,
107+
"file_id": "any-file-id",
108+
"status": "weird",
109+
})
110+
require.Error(t, err)
111+
})
112+
113+
t.Run("missing timestamp falls back to now without error", func(t *testing.T) {
114+
fs := inst.VFS()
115+
doc := createStatusTestFile(t, fs, "rag-no-ts.txt")
116+
defer destroyStatusTestFile(t, fs, doc)
117+
118+
err := runWorker(t, map[string]interface{}{
119+
"partition": inst.Domain,
120+
"file_id": doc.DocID,
121+
"status": "success",
122+
})
123+
require.NoError(t, err)
124+
125+
updated, err := fs.FileByID(doc.DocID)
126+
require.NoError(t, err)
127+
require.NotNil(t, updated.CozyMetadata.RAG.LastSuccessDate)
128+
})
129+
130+
t.Run("malformed timestamp falls back to now without error", func(t *testing.T) {
131+
fs := inst.VFS()
132+
doc := createStatusTestFile(t, fs, "rag-bad-ts.txt")
133+
defer destroyStatusTestFile(t, fs, doc)
134+
135+
err := runWorker(t, map[string]interface{}{
136+
"partition": inst.Domain,
137+
"file_id": doc.DocID,
138+
"status": "success",
139+
"timestamp": "not-a-date",
140+
})
141+
require.NoError(t, err)
142+
143+
updated, err := fs.FileByID(doc.DocID)
144+
require.NoError(t, err)
145+
require.NotNil(t, updated.CozyMetadata.RAG.LastSuccessDate)
146+
})
147+
148+
t.Run("non-existent file_id returns nil error", func(t *testing.T) {
149+
err := runWorker(t, map[string]interface{}{
150+
"partition": inst.Domain,
151+
"file_id": "non-existent-file-id",
152+
"status": "success",
153+
})
154+
require.NoError(t, err)
155+
})
156+
}
157+
158+
func createStatusTestFile(t *testing.T, fs vfs.VFS, name string) *vfs.FileDoc {
159+
t.Helper()
160+
parent, err := fs.DirByPath("/")
161+
require.NoError(t, err)
162+
doc, err := vfs.NewFileDoc(name, parent.DocID, 4, nil, "text/plain", "text", time.Now(), false, false, false, nil)
163+
require.NoError(t, err)
164+
f, err := fs.CreateFile(doc, nil)
165+
require.NoError(t, err)
166+
_, err = f.Write([]byte("test"))
167+
require.NoError(t, err)
168+
require.NoError(t, f.Close())
169+
updated, err := fs.FileByID(doc.DocID)
170+
require.NoError(t, err)
171+
return updated
172+
}
173+
174+
func destroyStatusTestFile(t *testing.T, fs vfs.VFS, doc *vfs.FileDoc) {
175+
t.Helper()
176+
_ = fs.DestroyFile(doc)
177+
}

0 commit comments

Comments
 (0)