diff --git a/internal/mountsync/syncer.go b/internal/mountsync/syncer.go index 075c7c5d..e469e8d5 100644 --- a/internal/mountsync/syncer.go +++ b/internal/mountsync/syncer.go @@ -30,6 +30,7 @@ import ( "unicode/utf8" "github.com/agentworkforce/relayfile/internal/digest" + "github.com/agentworkforce/relayfile/internal/relayfile" "github.com/fsnotify/fsnotify" "nhooyr.io/websocket" "nhooyr.io/websocket/wsjson" @@ -63,6 +64,11 @@ func (e *SchemaValidationError) Is(target error) bool { const defaultBulkFlushThreshold = 256 +const ( + mountWritebackCreateDraftContentIdentityKind = "mount-writeback-create-draft" + mountWritebackCreateDraftContentIdentityTTLSeconds = 2592000 +) + // defaultFullPullEvery is the default cadence for the "trust but verify" // periodic full tree pull that runs from the incremental path. At 30s sync // intervals, 20 cycles is roughly every 10 minutes. This is the safety net @@ -1577,7 +1583,7 @@ func (s *Syncer) flushPendingBulkWrites(ctx context.Context, pending []pendingBu s.logf("writeback flush refused: cloud-error circuit breaker is open; %d file(s) remain pending", len(pending)) return nil } - for _, chunk := range chunkPendingBulkWrites(pending, maxWritebackBatchBytes()) { + for _, chunk := range chunkPendingBulkWrites(s.workspace, pending, maxWritebackBatchBytes()) { if err := s.flushPendingBulkWriteChunk(ctx, chunk, conflicted); err != nil { return err } @@ -1586,16 +1592,7 @@ func (s *Syncer) flushPendingBulkWrites(ctx context.Context, pending []pendingBu } func (s *Syncer) flushPendingBulkWriteChunk(ctx context.Context, pending []pendingBulkWrite, conflicted map[string]struct{}) error { - files := make([]BulkWriteFile, 0, len(pending)) - for _, pendingWrite := range pending { - files = append(files, BulkWriteFile{ - Path: pendingWrite.remotePath, - ContentType: pendingWrite.snapshot.ContentType, - Content: pendingWrite.snapshot.WireContent, - Encoding: pendingWrite.snapshot.Encoding, - }) - } - + files := bulkWriteFilesForPending(s.workspace, pending) response, err := s.client.WriteFilesBulk(ctx, s.workspace, files) if err != nil { s.recordCloudFailure(err) @@ -1637,20 +1634,36 @@ func (s *Syncer) flushPendingBulkWriteChunk(ctx context.Context, pending []pendi return firstErr } -func bulkWriteFilesForPending(pending []pendingBulkWrite) []BulkWriteFile { +func bulkWriteFilesForPending(workspaceID string, pending []pendingBulkWrite) []BulkWriteFile { files := make([]BulkWriteFile, 0, len(pending)) for _, pendingWrite := range pending { files = append(files, BulkWriteFile{ - Path: pendingWrite.remotePath, - ContentType: pendingWrite.snapshot.ContentType, - Content: pendingWrite.snapshot.WireContent, - Encoding: pendingWrite.snapshot.Encoding, + Path: pendingWrite.remotePath, + ContentType: pendingWrite.snapshot.ContentType, + Content: pendingWrite.snapshot.WireContent, + Encoding: pendingWrite.snapshot.Encoding, + ContentIdentity: mountWritebackCreateDraftContentIdentity(workspaceID, pendingWrite.remotePath, pendingWrite.snapshot.Hash), }) } return files } -func chunkPendingBulkWrites(pending []pendingBulkWrite, maxBytes int64) [][]pendingBulkWrite { +func mountWritebackCreateDraftContentIdentity(workspaceID, normalizedRemotePath, contentHash string) *ContentIdentity { + if !relayfile.IsDraftFilePath(normalizedRemotePath) { + return nil + } + return newMountWritebackCreateDraftContentIdentity(workspaceID, normalizedRemotePath, contentHash) +} + +func newMountWritebackCreateDraftContentIdentity(workspaceID, normalizedRemotePath, contentHash string) *ContentIdentity { + return &ContentIdentity{ + Kind: mountWritebackCreateDraftContentIdentityKind, + Key: fmt.Sprintf("%s:%s:%s", workspaceID, normalizedRemotePath, contentHash), + TTLSeconds: mountWritebackCreateDraftContentIdentityTTLSeconds, + } +} + +func chunkPendingBulkWrites(workspaceID string, pending []pendingBulkWrite, maxBytes int64) [][]pendingBulkWrite { if len(pending) == 0 { return nil } @@ -1661,7 +1674,7 @@ func chunkPendingBulkWrites(pending []pendingBulkWrite, maxBytes int64) [][]pend current := make([]pendingBulkWrite, 0, len(pending)) for _, item := range pending { candidate := append(append([]pendingBulkWrite(nil), current...), item) - if len(current) > 0 && bulkWriteRequestSize(bulkWriteFilesForPending(candidate)) > maxBytes { + if len(current) > 0 && bulkWriteRequestSize(bulkWriteFilesForPending(workspaceID, candidate)) > maxBytes { chunks = append(chunks, append([]pendingBulkWrite(nil), current...)) current = current[:0] } diff --git a/internal/mountsync/syncer_test.go b/internal/mountsync/syncer_test.go index a934bdc6..2c09f928 100644 --- a/internal/mountsync/syncer_test.go +++ b/internal/mountsync/syncer_test.go @@ -1812,6 +1812,183 @@ func TestBulkWrite_SingleCallForNFiles(t *testing.T) { } } +func TestBulkWrite_ContentIdentityGoldenVector(t *testing.T) { + const ( + workspaceID = "ws_test" + remotePath = "/slack/channels/C123/messages/messages 5ab77d67.json" + content = "{\"channel\":\"C123\",\"text\":\"hello writeback idempotency\"}\n" + contentHash = "751f9591557700f69b5ceefcdec7ead8563a10f0a712c501a5028699be021511" + key = "ws_test:/slack/channels/C123/messages/messages 5ab77d67.json:751f9591557700f69b5ceefcdec7ead8563a10f0a712c501a5028699be021511" + ) + + snapshot := newLocalSnapshot(remotePath, []byte(content)) + if snapshot.Hash != contentHash { + t.Fatalf("golden vector hash = %q, want %q", snapshot.Hash, contentHash) + } + identity := newMountWritebackCreateDraftContentIdentity(workspaceID, remotePath, snapshot.Hash) + if identity.Kind != mountWritebackCreateDraftContentIdentityKind { + t.Fatalf("content identity kind = %q, want %q", identity.Kind, mountWritebackCreateDraftContentIdentityKind) + } + if identity.Key != key { + t.Fatalf("content identity key = %q, want %q", identity.Key, key) + } + if identity.TTLSeconds != mountWritebackCreateDraftContentIdentityTTLSeconds { + t.Fatalf("content identity ttl = %d, want %d", identity.TTLSeconds, mountWritebackCreateDraftContentIdentityTTLSeconds) + } + if identity.TTLSeconds != 2592000 { + t.Fatalf("golden vector ttl = %d, want 2592000", identity.TTLSeconds) + } + if strings.TrimSpace(identity.Key) != identity.Key { + t.Fatalf("golden vector key should be trim-stable: %q", identity.Key) + } + if !strings.Contains(identity.Key, "messages 5ab77d67.json") { + t.Fatalf("golden vector key should preserve the internal path space: %q", identity.Key) + } +} + +func TestBulkWrite_ContentIdentityStabilityAndIsolation(t *testing.T) { + const ( + workspaceID = "ws_test" + remotePath = "/slack/channels/C123/messages/messages 5ab77d67-1111-4111-8111-123456789abc.json" + ) + + snapshot := newLocalSnapshot(remotePath, []byte("{\"text\":\"same\"}\n")) + files := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: remotePath, + snapshot: snapshot, + }}) + reupload := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: remotePath, + snapshot: snapshot, + }}) + if files[0].ContentIdentity == nil || reupload[0].ContentIdentity == nil { + t.Fatal("expected content identity on both bulk files") + } + if files[0].ContentIdentity.Key != reupload[0].ContentIdentity.Key { + t.Fatalf("same draft re-upload key changed: %q vs %q", files[0].ContentIdentity.Key, reupload[0].ContentIdentity.Key) + } + + edited := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: remotePath, + snapshot: newLocalSnapshot(remotePath, []byte("{\"text\":\"edited\"}\n")), + }}) + if edited[0].ContentIdentity.Key == files[0].ContentIdentity.Key { + t.Fatalf("edited draft content should change key %q", edited[0].ContentIdentity.Key) + } + + otherPath := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: "/slack/channels/C123/messages/messages 6ab77d67-2222-4222-8222-123456789abc.json", + snapshot: snapshot, + }}) + if otherPath[0].ContentIdentity.Key == files[0].ContentIdentity.Key { + t.Fatalf("different draft path should change key %q", otherPath[0].ContentIdentity.Key) + } + + otherWorkspace := bulkWriteFilesForPending("ws_other", []pendingBulkWrite{{ + remotePath: remotePath, + snapshot: snapshot, + }}) + if otherWorkspace[0].ContentIdentity.Key == files[0].ContentIdentity.Key { + t.Fatalf("different workspace should change key %q", otherWorkspace[0].ContentIdentity.Key) + } +} + +func TestBulkWrite_ContentIdentityOnlyForCreateDraftPaths(t *testing.T) { + const workspaceID = "ws_test" + snapshot := newLocalSnapshot("/notion/pages/pages 5ab77d67-1111-4111-8111-123456789abc.json", []byte("{}\n")) + + draft := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: "/notion/pages/pages 5ab77d67-1111-4111-8111-123456789abc.json", + snapshot: snapshot, + }}) + if draft[0].ContentIdentity == nil { + t.Fatal("expected content identity for space-uuid create draft path") + } + + stable := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: "/notion/pages/page-1.md", + snapshot: snapshot, + }}) + if stable[0].ContentIdentity != nil { + t.Fatalf("stable non-draft path should not carry content identity: %+v", stable[0].ContentIdentity) + } + + nonUUID := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: "/slack/channels/C123/messages/messages not-a-uuid.json", + snapshot: snapshot, + }}) + if nonUUID[0].ContentIdentity != nil { + t.Fatalf("non-uuid draft-like path should not carry content identity: %+v", nonUUID[0].ContentIdentity) + } +} + +func TestBulkWrite_ContentIdentityOmittedForStablePathRevert(t *testing.T) { + const ( + workspaceID = "ws_test" + remotePath = "/notion/pages/page-1.md" + ) + + for _, content := range []string{"# C\n", "# D\n", "# C\n"} { + files := bulkWriteFilesForPending(workspaceID, []pendingBulkWrite{{ + remotePath: remotePath, + snapshot: newLocalSnapshot(remotePath, []byte(content)), + }}) + if files[0].ContentIdentity != nil { + t.Fatalf("stable path content %q should not carry content identity: %+v", content, files[0].ContentIdentity) + } + } +} + +func TestBulkWrite_FlushSendsContentIdentity(t *testing.T) { + const ( + workspaceID = "ws_test" + content = "{\"channel\":\"C123\",\"text\":\"hello writeback idempotency\"}\n" + draftID = "5ab77d67-1111-4111-8111-123456789abc" + key = "ws_test:/slack/channels/C123/messages/messages 5ab77d67-1111-4111-8111-123456789abc.json:751f9591557700f69b5ceefcdec7ead8563a10f0a712c501a5028699be021511" + ) + + client := &fakeClient{files: map[string]RemoteFile{}} + localDir := t.TempDir() + messageDir := filepath.Join(localDir, "slack", "channels", "C123", "messages") + if err := os.MkdirAll(messageDir, 0o755); err != nil { + t.Fatalf("mkdir message dir failed: %v", err) + } + if err := os.WriteFile(filepath.Join(messageDir, "messages "+draftID+".json"), []byte(content), 0o644); err != nil { + t.Fatalf("write draft failed: %v", err) + } + + syncer, err := NewSyncer(client, SyncerOptions{ + WorkspaceID: workspaceID, + RemoteRoot: "/", + LocalRoot: localDir, + }) + if err != nil { + t.Fatalf("new syncer failed: %v", err) + } + if err := syncer.SyncOnce(context.Background()); err != nil { + t.Fatalf("sync once failed: %v", err) + } + if got := len(client.bulkWriteBatches); got != 1 { + t.Fatalf("expected one bulk write batch, got %d", got) + } + if got := len(client.bulkWriteBatches[0]); got != 1 { + t.Fatalf("expected one bulk write file, got %d", got) + } + identity := client.bulkWriteBatches[0][0].ContentIdentity + if identity == nil { + t.Fatal("expected flushed bulk file to carry content identity") + } + if identity.Kind != mountWritebackCreateDraftContentIdentityKind { + t.Fatalf("flushed identity kind = %q, want %q", identity.Kind, mountWritebackCreateDraftContentIdentityKind) + } + if identity.Key != key { + t.Fatalf("flushed identity key = %q, want %q", identity.Key, key) + } + if identity.TTLSeconds != 2592000 { + t.Fatalf("flushed identity ttl = %d, want 2592000", identity.TTLSeconds) + } +} + func TestBulkWrite_MixedCreateAndUpdateBatch(t *testing.T) { client := &fakeClient{ files: map[string]RemoteFile{ diff --git a/internal/mountsync/types.go b/internal/mountsync/types.go index 32c882f2..a2035f9c 100644 --- a/internal/mountsync/types.go +++ b/internal/mountsync/types.go @@ -9,6 +9,8 @@ import ( type BulkWriteFile = relayfile.BulkWriteFile +type ContentIdentity = relayfile.ContentIdentity + type BulkWriteError = relayfile.BulkWriteError type BulkWriteResult = relayfile.BulkWriteResult diff --git a/internal/relayfile/draft_reconcile.go b/internal/relayfile/draft_reconcile.go index f16f2dde..cfea7437 100644 --- a/internal/relayfile/draft_reconcile.go +++ b/internal/relayfile/draft_reconcile.go @@ -104,6 +104,12 @@ func isDraftFileBasename(basename string) bool { return draftFileBasenamePattern.MatchString(basename) } +// IsDraftFilePath reports whether a path matches the relayfile-owned +// draftFile() create-draft basename contract. +func IsDraftFilePath(filePath string) bool { + return isDraftFileBasename(basenameOf(filePath)) +} + // reconcileAckedDraftLocked applies the draftFile() rename contract after a // successful externally-executed writeback: the agent-authored draft is // renamed to the canonical id (or removed when the canonical record already diff --git a/internal/relayfile/store.go b/internal/relayfile/store.go index 1a9f1596..2d2af51c 100644 --- a/internal/relayfile/store.go +++ b/internal/relayfile/store.go @@ -171,11 +171,29 @@ type WriteRequest struct { CorrelationID string } +type ContentIdentity struct { + Kind string `json:"kind"` + Key string `json:"key"` + TTLSeconds int `json:"ttlSeconds,omitempty"` +} + +func cloneContentIdentity(identity *ContentIdentity) *ContentIdentity { + if identity == nil { + return nil + } + return &ContentIdentity{ + Kind: identity.Kind, + Key: identity.Key, + TTLSeconds: identity.TTLSeconds, + } +} + type BulkWriteFile struct { - Path string `json:"path"` - ContentType string `json:"contentType"` - Content string `json:"content"` - Encoding string `json:"encoding"` + Path string `json:"path"` + ContentType string `json:"contentType"` + Content string `json:"content"` + Encoding string `json:"encoding"` + ContentIdentity *ContentIdentity `json:"contentIdentity,omitempty"` } type BulkWriteError struct { @@ -208,20 +226,21 @@ type WriteResult struct { } type OperationStatus struct { - OpID string `json:"opId"` - Path string `json:"path,omitempty"` - Revision string `json:"revision,omitempty"` - Action string `json:"action,omitempty"` - Provider string `json:"provider,omitempty"` - Status string `json:"status"` - AttemptCount int `json:"attemptCount"` - NextAttemptAt *string `json:"nextAttemptAt,omitempty"` - LastError *string `json:"lastError,omitempty"` - ProviderResult map[string]any `json:"providerResult,omitempty"` - CorrelationID string `json:"correlationId,omitempty"` - CreatedAt string `json:"createdAt,omitempty"` - UpdatedAt string `json:"updatedAt,omitempty"` - CompletedAt *string `json:"completedAt,omitempty"` + OpID string `json:"opId"` + Path string `json:"path,omitempty"` + Revision string `json:"revision,omitempty"` + Action string `json:"action,omitempty"` + Provider string `json:"provider,omitempty"` + Status string `json:"status"` + AttemptCount int `json:"attemptCount"` + NextAttemptAt *string `json:"nextAttemptAt,omitempty"` + LastError *string `json:"lastError,omitempty"` + ProviderResult map[string]any `json:"providerResult,omitempty"` + ContentIdentity *ContentIdentity `json:"contentIdentity,omitempty"` + CorrelationID string `json:"correlationId,omitempty"` + CreatedAt string `json:"createdAt,omitempty"` + UpdatedAt string `json:"updatedAt,omitempty"` + CompletedAt *string `json:"completedAt,omitempty"` } type OperationFeed struct { @@ -392,6 +411,7 @@ type WritebackAction struct { Type WritebackActionType `json:"type"` ContentType string `json:"contentType,omitempty"` Content string `json:"content,omitempty"` + ContentIdentity *ContentIdentity `json:"contentIdentity,omitempty"` Provider string `json:"provider,omitempty"` ProviderObjectID string `json:"providerObjectId,omitempty"` CorrelationID string `json:"correlationId,omitempty"` @@ -512,11 +532,12 @@ type workspaceState struct { } type WritebackQueueItem struct { - WorkspaceID string `json:"workspaceId"` - OpID string `json:"opId"` - Path string `json:"path"` - Revision string `json:"revision"` - CorrelationID string `json:"correlationId"` + WorkspaceID string `json:"workspaceId"` + OpID string `json:"opId"` + Path string `json:"path"` + Revision string `json:"revision"` + ContentIdentity *ContentIdentity `json:"contentIdentity,omitempty"` + CorrelationID string `json:"correlationId"` } type writebackTask = WritebackQueueItem @@ -967,11 +988,12 @@ func NewStoreWithOptions(opts StoreOptions) *Store { continue } task := writebackTask{ - WorkspaceID: workspaceID, - OpID: opID, - Path: op.Path, - Revision: op.Revision, - CorrelationID: op.CorrelationID, + WorkspaceID: workspaceID, + OpID: opID, + Path: op.Path, + Revision: op.Revision, + ContentIdentity: cloneContentIdentity(op.ContentIdentity), + CorrelationID: op.CorrelationID, } delay := time.Duration(0) if op.NextAttemptAt != nil { @@ -1443,7 +1465,7 @@ func (s *Store) BulkWrite(workspaceID string, files []BulkWriteFile) (int, []Bul if existed { eventType = "file.updated" } - result, task := s.recordWriteLocked(ws, path, revision, eventType, file.Provider, "") + result, task := s.recordWriteWithContentIdentityLocked(ws, path, revision, eventType, file.Provider, "", input.ContentIdentity) _ = result results = append(results, BulkWriteResult{ Path: path, @@ -3070,13 +3092,19 @@ func (s *Store) GetPendingWritebacks(workspaceID string) []map[string]any { if op.Status != "pending" && op.Status != "running" { continue } - result = append(result, map[string]any{ + itemOut := map[string]any{ "id": item.OpID, "workspaceId": item.WorkspaceID, "path": item.Path, "revision": item.Revision, "correlationId": item.CorrelationID, - }) + } + if item.ContentIdentity != nil { + itemOut["contentIdentity"] = item.ContentIdentity + } else if op.ContentIdentity != nil { + itemOut["contentIdentity"] = op.ContentIdentity + } + result = append(result, itemOut) } return result } @@ -3420,22 +3448,27 @@ func nowRFC3339NanoUTC() string { } func (s *Store) recordWriteLocked(ws *workspaceState, path, revision, eventType, provider, correlationID string) (WriteResult, writebackTask) { + return s.recordWriteWithContentIdentityLocked(ws, path, revision, eventType, provider, correlationID, nil) +} + +func (s *Store) recordWriteWithContentIdentityLocked(ws *workspaceState, path, revision, eventType, provider, correlationID string, contentIdentity *ContentIdentity) (WriteResult, writebackTask) { if provider == "" { } workspaceID := s.workspaceIDForStateLocked(ws) opID := s.nextOperationIDLocked() nowTS := nowRFC3339NanoUTC() op := OperationStatus{ - OpID: opID, - Path: path, - Revision: revision, - Action: string(writebackActionFromEventType(eventType)), - Provider: provider, - Status: "pending", - AttemptCount: 0, - CorrelationID: correlationID, - CreatedAt: nowTS, - UpdatedAt: nowTS, + OpID: opID, + Path: path, + Revision: revision, + Action: string(writebackActionFromEventType(eventType)), + Provider: provider, + Status: "pending", + AttemptCount: 0, + ContentIdentity: cloneContentIdentity(contentIdentity), + CorrelationID: correlationID, + CreatedAt: nowTS, + UpdatedAt: nowTS, } ws.Ops[opID] = op @@ -3460,7 +3493,7 @@ func (s *Store) recordWriteLocked(ws *workspaceState, path, revision, eventType, result.Writeback.Provider = provider result.Writeback.State = "pending" - task := writebackTask{WorkspaceID: workspaceID, OpID: opID, Path: path, Revision: revision, CorrelationID: correlationID} + task := writebackTask{WorkspaceID: workspaceID, OpID: opID, Path: path, Revision: revision, ContentIdentity: cloneContentIdentity(contentIdentity), CorrelationID: correlationID} return result, task } @@ -3806,10 +3839,14 @@ func (s *Store) processWriteback(task writebackTask) { task.Revision = op.Revision } writeAction := WritebackAction{ - WorkspaceID: task.WorkspaceID, - Path: task.Path, - Revision: task.Revision, - CorrelationID: task.CorrelationID, + WorkspaceID: task.WorkspaceID, + Path: task.Path, + Revision: task.Revision, + ContentIdentity: cloneContentIdentity(task.ContentIdentity), + CorrelationID: task.CorrelationID, + } + if writeAction.ContentIdentity == nil { + writeAction.ContentIdentity = cloneContentIdentity(op.ContentIdentity) } if op.Provider != "" { writeAction.Provider = op.Provider diff --git a/internal/relayfile/store_test.go b/internal/relayfile/store_test.go index 44cd1a2c..f74c4ca7 100644 --- a/internal/relayfile/store_test.go +++ b/internal/relayfile/store_test.go @@ -4272,6 +4272,47 @@ func TestProviderWriteActionReceivesFileUpsertPayload(t *testing.T) { } } +func TestBulkWriteContentIdentityReachesProviderWriteAction(t *testing.T) { + actions := make(chan WritebackAction, 1) + store := NewStoreWithOptions(StoreOptions{ + ProviderWriteAction: func(action WritebackAction) error { + actions <- action + return nil + }, + }) + t.Cleanup(store.Close) + + identity := &ContentIdentity{ + Kind: "mount-writeback-create-draft", + Key: "ws_bulk_identity:/external/Draft.md:abc123", + TTLSeconds: 2592000, + } + written, _, errs := store.BulkWrite("ws_bulk_identity", []BulkWriteFile{{ + Path: "/external/Draft.md", + ContentType: "text/markdown", + Content: "# draft", + ContentIdentity: identity, + }}) + if len(errs) != 0 { + t.Fatalf("bulk write returned errors: %+v", errs) + } + if written != 1 { + t.Fatalf("expected one bulk write, got %d", written) + } + + select { + case action := <-actions: + if action.ContentIdentity == nil { + t.Fatal("expected content identity on provider write action") + } + if *action.ContentIdentity != *identity { + t.Fatalf("provider write content identity = %+v, want %+v", action.ContentIdentity, identity) + } + case <-time.After(2 * time.Second): + t.Fatalf("expected provider write action callback") + } +} + func TestProviderWriteActionReceivesFileDeletePayload(t *testing.T) { actions := make(chan WritebackAction, 10) store := NewStoreWithOptions(StoreOptions{ @@ -4746,3 +4787,40 @@ func TestExternalWritebackModeKeepsItemsInQueue(t *testing.T) { } } } + +func TestBulkWriteContentIdentityAppearsInPendingWritebacks(t *testing.T) { + store := NewStoreWithOptions(StoreOptions{ + ExternalWritebackMode: true, + }) + t.Cleanup(store.Close) + + identity := &ContentIdentity{ + Kind: "mount-writeback-create-draft", + Key: "ws_ext_identity:/external/Draft.md:abc123", + TTLSeconds: 2592000, + } + written, _, errs := store.BulkWrite("ws_ext_identity", []BulkWriteFile{{ + Path: "/external/Draft.md", + ContentType: "text/markdown", + Content: "# external draft", + ContentIdentity: identity, + }}) + if len(errs) != 0 { + t.Fatalf("bulk write returned errors: %+v", errs) + } + if written != 1 { + t.Fatalf("expected one bulk write, got %d", written) + } + + pending := store.GetPendingWritebacks("ws_ext_identity") + if len(pending) != 1 { + t.Fatalf("expected one pending writeback, got %+v", pending) + } + rawIdentity, ok := pending[0]["contentIdentity"].(*ContentIdentity) + if !ok || rawIdentity == nil { + t.Fatalf("expected pending content identity, got %+v", pending[0]["contentIdentity"]) + } + if *rawIdentity != *identity { + t.Fatalf("pending content identity = %+v, want %+v", rawIdentity, identity) + } +} diff --git a/openapi/relayfile-v1.openapi.yaml b/openapi/relayfile-v1.openapi.yaml index 39591a3c..a77d60bd 100644 --- a/openapi/relayfile-v1.openapi.yaml +++ b/openapi/relayfile-v1.openapi.yaml @@ -2460,6 +2460,8 @@ components: providerResult: type: object additionalProperties: true + contentIdentity: + $ref: '#/components/schemas/ContentIdentity' correlationId: type: string createdAt: @@ -2847,6 +2849,8 @@ components: correlationId: type: string description: Correlation ID for tracing + contentIdentity: + $ref: '#/components/schemas/ContentIdentity' WritebackItem: type: object @@ -2952,6 +2956,26 @@ components: encoding: type: string description: Content encoding, e.g. "base64" for binary data + contentIdentity: + $ref: '#/components/schemas/ContentIdentity' + + ContentIdentity: + type: object + additionalProperties: false + required: [kind, key] + properties: + kind: + type: string + minLength: 1 + description: Dedupe namespace for this write identity + key: + type: string + minLength: 1 + description: Caller-provided idempotency key within the namespace + ttlSeconds: + type: integer + minimum: 0 + description: Optional lifetime for the dedupe identity in seconds BulkWriteError: type: object diff --git a/packages/core/src/dedup.coverage.test.ts b/packages/core/src/dedup.coverage.test.ts index a2010518..047370fd 100644 --- a/packages/core/src/dedup.coverage.test.ts +++ b/packages/core/src/dedup.coverage.test.ts @@ -23,6 +23,7 @@ describe("dedup type contracts", () => { const identity: ContentIdentity = { kind: "push", key: "repo:abc:main", + ttlSeconds: 300, }; const store: DedupStore = { async has(key) { diff --git a/packages/core/src/dedup.ts b/packages/core/src/dedup.ts index eccab355..5f3f3c9e 100644 --- a/packages/core/src/dedup.ts +++ b/packages/core/src/dedup.ts @@ -3,6 +3,7 @@ import { createHash } from "node:crypto"; export interface ContentIdentity { kind: string; key: string; + ttlSeconds?: number; } export interface DedupEntry { diff --git a/packages/sdk/typescript/src/client.test.ts b/packages/sdk/typescript/src/client.test.ts index be8a9972..b807db7f 100644 --- a/packages/sdk/typescript/src/client.test.ts +++ b/packages/sdk/typescript/src/client.test.ts @@ -1176,10 +1176,18 @@ describe("RelayFileClient — existing methods", () => { path: "/github/push/abc123.json", baseRevision: "rev_3", content: "{}", - contentIdentity: { kind: "github.push", key: "abc123" }, + contentIdentity: { + kind: "github.push", + key: "abc123", + ttlSeconds: 2592000, + }, }); const body = JSON.parse((f.mock.calls[0]![1] as RequestInit).body as string); - expect(body.contentIdentity).toEqual({ kind: "github.push", key: "abc123" }); + expect(body.contentIdentity).toEqual({ + kind: "github.push", + key: "abc123", + ttlSeconds: 2592000, + }); }); it("omits contentIdentity from the body when not provided", async () => { @@ -1270,7 +1278,15 @@ describe("RelayFileClient — existing methods", () => { const res: BulkWriteResponse = await client.bulkWrite({ workspaceId: "ws_acme", files: [ - { path: "/a.md", content: "a" }, + { + path: "/a.md", + content: "a", + contentIdentity: { + kind: "mount-writeback-create-draft", + key: "ws_acme:/a.md:hash", + ttlSeconds: 2592000, + }, + }, { path: "/b.md", content: "b", encoding: "utf-8" }, ], }); @@ -1292,7 +1308,15 @@ describe("RelayFileClient — existing methods", () => { expect(init.method).toBe("POST"); expect(JSON.parse(init.body as string)).toEqual({ files: [ - { path: "/a.md", content: "a" }, + { + path: "/a.md", + content: "a", + contentIdentity: { + kind: "mount-writeback-create-draft", + key: "ws_acme:/a.md:hash", + ttlSeconds: 2592000, + }, + }, { path: "/b.md", content: "b", encoding: "utf-8" }, ], }); diff --git a/packages/sdk/typescript/src/types.ts b/packages/sdk/typescript/src/types.ts index 94d27252..83d040af 100644 --- a/packages/sdk/typescript/src/types.ts +++ b/packages/sdk/typescript/src/types.ts @@ -59,6 +59,7 @@ export interface FileSemantics { export interface ContentIdentity { kind: string; key: string; + ttlSeconds?: number; } export interface FileReadResponse {