Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 31 additions & 18 deletions internal/mountsync/syncer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
}
Expand All @@ -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)
Expand Down Expand Up @@ -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
}
Expand All @@ -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]
}
Expand Down
177 changes: 177 additions & 0 deletions internal/mountsync/syncer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Comment on lines +1829 to +1831

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Pin the ContentIdentity.Kind wire literal in these contract tests

Both assertions currently compare against an internal constant, so a constant drift could pass tests while still breaking the locked wire contract. Assert the literal "mount-writeback-create-draft" in the test path.

🔧 Suggested patch
@@
-	if identity.Kind != mountWritebackCreateDraftContentIdentityKind {
-		t.Fatalf("content identity kind = %q, want %q", identity.Kind, mountWritebackCreateDraftContentIdentityKind)
+	if identity.Kind != "mount-writeback-create-draft" {
+		t.Fatalf("content identity kind = %q, want %q", identity.Kind, "mount-writeback-create-draft")
 	}
@@
-	if identity.Kind != mountWritebackCreateDraftContentIdentityKind {
-		t.Fatalf("flushed identity kind = %q, want %q", identity.Kind, mountWritebackCreateDraftContentIdentityKind)
+	if identity.Kind != "mount-writeback-create-draft" {
+		t.Fatalf("flushed identity kind = %q, want %q", identity.Kind, "mount-writeback-create-draft")
 	}

Also applies to: 1947-1949

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@internal/mountsync/syncer_test.go` around lines 1842 - 1844, The test
currently asserts identity.Kind against the internal constant
mountWritebackCreateDraftContentIdentityKind which can mask contract drift;
update the assertions in internal/mountsync/syncer_test.go to compare
identity.Kind to the literal string "mount-writeback-create-draft" instead of
the constant (replace uses of mountWritebackCreateDraftContentIdentityKind with
the wire literal) in both places around the existing identity.Kind checks so the
test verifies the actual on-the-wire value.

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{
Expand Down
2 changes: 2 additions & 0 deletions internal/mountsync/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ import (

type BulkWriteFile = relayfile.BulkWriteFile

type ContentIdentity = relayfile.ContentIdentity

type BulkWriteError = relayfile.BulkWriteError

type BulkWriteResult = relayfile.BulkWriteResult
Expand Down
6 changes: 6 additions & 0 deletions internal/relayfile/draft_reconcile.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading