From 012c258ceee86c12db6d4ef8a024b6fc7ea1b7a2 Mon Sep 17 00:00:00 2001 From: Otavio Carvalho Date: Fri, 2 Oct 2026 10:17:30 +0000 Subject: [PATCH] Add append mode to WriteFile WriteFile currently opens the target with O_TRUNC on every stream, so guest-side transcript files cannot be grown across writes. Add a bool append field to WriteFileRequest (proto field 4): when set on the first message the guest opens O_CREATE|O_WRONLY|O_APPEND instead of O_CREATE|O_WRONLY|O_TRUNC; the file is created if it does not exist. The api proxies WriteFileRequest opaquely, so no api-side change is needed. --- guest/filesystem/service.go | 8 +- guest/filesystem/service_test.go | 162 +++++++++++++++++++++++++++++ internal/apiservice/server_test.go | 82 +++++++++++++++ proto/ateenv/v1alpha/guest.pb.go | 17 ++- proto/ateenv/v1alpha/guest.proto | 3 + 5 files changed, 268 insertions(+), 4 deletions(-) diff --git a/guest/filesystem/service.go b/guest/filesystem/service.go index 5162969..a698b46 100644 --- a/guest/filesystem/service.go +++ b/guest/filesystem/service.go @@ -207,7 +207,13 @@ func (s *Service) WriteFile(stream ateenvv1alpha.FileSystemService_WriteFileServ mode = os.FileMode(req.GetMode()) } - f, err = os.OpenFile(filePath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, mode) + // Append streams keep existing content; default streams truncate. + flags := os.O_CREATE | os.O_WRONLY | os.O_TRUNC + if req.GetAppend() { + flags = os.O_CREATE | os.O_WRONLY | os.O_APPEND + } + + f, err = os.OpenFile(filePath, flags, mode) if err != nil { if errors.Is(err, os.ErrPermission) { return status.Errorf(codes.PermissionDenied, "permission denied opening %q: %v", reqPath, err) diff --git a/guest/filesystem/service_test.go b/guest/filesystem/service_test.go index bc78615..172c499 100644 --- a/guest/filesystem/service_test.go +++ b/guest/filesystem/service_test.go @@ -349,3 +349,165 @@ func TestWriteFileMissingPath(t *testing.T) { t.Fatalf("expected InvalidArgument for missing path, got %v", err) } } + +func TestAppendPreservesExistingContent(t *testing.T) { + tempDir := t.TempDir() + client, cleanup := setupTestFileSystemServer(t, Config{RootDirectory: tempDir, ReadBufferSize: 4 * 1024}) + defer cleanup() + + ctx := context.Background() + targetPath := filepath.Join(tempDir, "transcript.jsonl") + initial := []byte("{\"seq\":0}\n") + second := []byte("{\"seq\":1}\n") + third := []byte("{\"seq\":2}\n") + + // Seed the file with a default (truncating) write. + writeStream, err := client.WriteFile(ctx) + if err != nil { + t.Fatalf("WriteFile failed: %v", err) + } + if err := writeStream.Send(&ateenvv1alpha.WriteFileRequest{ + Path: targetPath, + Chunk: initial, + }); err != nil { + t.Fatalf("failed to send seed chunk: %v", err) + } + if _, err := writeStream.CloseAndRecv(); err != nil { + t.Fatalf("seed write CloseAndRecv failed: %v", err) + } + + // Append two more chunks with append=true on the first message only. + appendStream, err := client.WriteFile(ctx) + if err != nil { + t.Fatalf("WriteFile failed: %v", err) + } + if err := appendStream.Send(&ateenvv1alpha.WriteFileRequest{ + Path: targetPath, + Chunk: second, + Append: true, + }); err != nil { + t.Fatalf("failed to send append chunk: %v", err) + } + if err := appendStream.Send(&ateenvv1alpha.WriteFileRequest{ + Chunk: third, + }); err != nil { + t.Fatalf("failed to send second append chunk: %v", err) + } + appendRes, err := appendStream.CloseAndRecv() + if err != nil { + t.Fatalf("append CloseAndRecv failed: %v", err) + } + if got, want := appendRes.BytesWritten, int64(len(second)+len(third)); got != want { + t.Fatalf("append bytes written = %d, want %d", got, want) + } + + // Read back: original content preserved, appended chunks in order. + got := readFileContent(t, client, ctx, targetPath) + want := string(initial) + string(second) + string(third) + if got != want { + t.Fatalf("content after append = %q, want %q", got, want) + } +} + +func TestAppendCreatesNewFile(t *testing.T) { + tempDir := t.TempDir() + client, cleanup := setupTestFileSystemServer(t, Config{RootDirectory: tempDir, ReadBufferSize: 4 * 1024}) + defer cleanup() + + ctx := context.Background() + targetPath := filepath.Join(tempDir, "brand-new.log") + content := []byte("first line\n") + + writeStream, err := client.WriteFile(ctx) + if err != nil { + t.Fatalf("WriteFile failed: %v", err) + } + if err := writeStream.Send(&ateenvv1alpha.WriteFileRequest{ + Path: targetPath, + Chunk: content, + Mode: 0640, + Append: true, + }); err != nil { + t.Fatalf("failed to send chunk: %v", err) + } + res, err := writeStream.CloseAndRecv() + if err != nil { + t.Fatalf("append-to-new-file CloseAndRecv failed: %v", err) + } + if res.BytesWritten != int64(len(content)) { + t.Fatalf("bytes written = %d, want %d", res.BytesWritten, len(content)) + } + + info, err := os.Stat(targetPath) + if err != nil { + t.Fatalf("failed to stat written file: %v", err) + } + if info.Mode().Perm() != 0640 { + t.Fatalf("expected permissions 0640, got %v", info.Mode().Perm()) + } + if got := readFileContent(t, client, ctx, targetPath); got != string(content) { + t.Fatalf("content = %q, want %q", got, string(content)) + } +} + +func TestDefaultWriteStillTruncates(t *testing.T) { + tempDir := t.TempDir() + client, cleanup := setupTestFileSystemServer(t, Config{RootDirectory: tempDir, ReadBufferSize: 4 * 1024}) + defer cleanup() + + ctx := context.Background() + targetPath := filepath.Join(tempDir, "replaced.txt") + + writeStream, err := client.WriteFile(ctx) + if err != nil { + t.Fatalf("WriteFile failed: %v", err) + } + if err := writeStream.Send(&ateenvv1alpha.WriteFileRequest{ + Path: targetPath, + Chunk: []byte("original content that must go away\n"), + }); err != nil { + t.Fatalf("failed to send seed chunk: %v", err) + } + if _, err := writeStream.CloseAndRecv(); err != nil { + t.Fatalf("seed write CloseAndRecv failed: %v", err) + } + + writeStream, err = client.WriteFile(ctx) + if err != nil { + t.Fatalf("WriteFile failed: %v", err) + } + if err := writeStream.Send(&ateenvv1alpha.WriteFileRequest{ + Path: targetPath, + Chunk: []byte("new"), + }); err != nil { + t.Fatalf("failed to send chunk: %v", err) + } + if _, err := writeStream.CloseAndRecv(); err != nil { + t.Fatalf("default write CloseAndRecv failed: %v", err) + } + + if got, want := readFileContent(t, client, ctx, targetPath), "new"; got != want { + t.Fatalf("content after default write = %q, want %q", got, want) + } +} + +// readFileContent streams a file back and returns its full content. +func readFileContent(t *testing.T, client ateenvv1alpha.FileSystemServiceClient, ctx context.Context, path string) string { + t.Helper() + readStream, err := client.ReadFile(ctx, &ateenvv1alpha.ReadFileRequest{Path: path}) + if err != nil { + t.Fatalf("ReadFile failed: %v", err) + } + var buf bytes.Buffer + for { + chunk, err := readStream.Recv() + if err == io.EOF { + break + } + if err != nil { + t.Fatalf("ReadFile recv error: %v", err) + } + buf.Write(chunk.Data) + } + return buf.String() +} diff --git a/internal/apiservice/server_test.go b/internal/apiservice/server_test.go index 56578a9..e9b2362 100644 --- a/internal/apiservice/server_test.go +++ b/internal/apiservice/server_test.go @@ -431,6 +431,88 @@ func TestProxyGuestServices(t *testing.T) { } } +// TestProxyAppendWrite verifies the WriteFileRequest append flag survives the +// api proxy and appends at the guest. +func TestProxyAppendWrite(t *testing.T) { + te := newFullTestEnv(t) + ctx := context.Background() + + if _, err := te.envClient.CreateEnvironment(ctx, &ateenvv1alpha.CreateEnvironmentRequest{ + Id: "append-test", + }); err != nil { + t.Fatalf("CreateEnvironment: %v", err) + } + workDir := t.TempDir() + grpcGuestServer, cleanup, err := guest.NewServer(guest.Config{ + Workspace: workDir, + EnableProcess: true, + EnableFileSystem: true, + }) + if err != nil { + t.Fatalf("guest.NewServer: %v", err) + } + t.Cleanup(cleanup) + te.router.Register("append-test", grpcGuestServer) + + envCtx := metadata.AppendToOutgoingContext(ctx, "x-env-id", "append-test", "x-env-atespace", "default") + + // Default write seeds the file. + writeStream, err := te.fsClient.WriteFile(envCtx) + if err != nil { + t.Fatalf("WriteFile stream: %v", err) + } + if err := writeStream.Send(&ateenvv1alpha.WriteFileRequest{ + Path: "append-file.txt", + Chunk: []byte("line-1\n"), + }); err != nil { + t.Fatalf("WriteFile send: %v", err) + } + if _, err := writeStream.CloseAndRecv(); err != nil { + t.Fatalf("WriteFile CloseAndRecv: %v", err) + } + + // Append through the proxy with the flag set on the first message. + appendStream, err := te.fsClient.WriteFile(envCtx) + if err != nil { + t.Fatalf("append WriteFile stream: %v", err) + } + if err := appendStream.Send(&ateenvv1alpha.WriteFileRequest{ + Path: "append-file.txt", + Chunk: []byte("line-2\n"), + Append: true, + }); err != nil { + t.Fatalf("append WriteFile send: %v", err) + } + resp, err := appendStream.CloseAndRecv() + if err != nil { + t.Fatalf("append WriteFile CloseAndRecv: %v", err) + } + if resp.GetBytesWritten() != int64(len("line-2\n")) { + t.Errorf("append bytes written = %d, want %d", resp.GetBytesWritten(), len("line-2\n")) + } + + readStream, err := te.fsClient.ReadFile(envCtx, &ateenvv1alpha.ReadFileRequest{ + Path: "append-file.txt", + }) + if err != nil { + t.Fatalf("ReadFile stream: %v", err) + } + var readBuf []byte + for { + chunk, err := readStream.Recv() + if err == io.EOF { + break + } + if err != nil { + t.Fatalf("ReadFile recv: %v", err) + } + readBuf = append(readBuf, chunk.GetData()...) + } + if string(readBuf) != "line-1\nline-2\n" { + t.Errorf("read %q, want %q", string(readBuf), "line-1\nline-2\n") + } +} + func TestActorStatusToEnvStatus(t *testing.T) { cases := []struct { name string diff --git a/proto/ateenv/v1alpha/guest.pb.go b/proto/ateenv/v1alpha/guest.pb.go index 53a1cb7..c28a92d 100644 --- a/proto/ateenv/v1alpha/guest.pb.go +++ b/proto/ateenv/v1alpha/guest.pb.go @@ -711,7 +711,10 @@ type WriteFileRequest struct { // Raw binary or text chunk to write to the file. Chunk []byte `protobuf:"bytes,2,opt,name=chunk,proto3" json:"chunk,omitempty"` // Optional POSIX file mode permission (e.g. 0644 or 0755; processed on first message). - Mode uint32 `protobuf:"varint,3,opt,name=mode,proto3" json:"mode,omitempty"` + Mode uint32 `protobuf:"varint,3,opt,name=mode,proto3" json:"mode,omitempty"` + // Append to the file instead of truncating it. The file is created if it + // does not exist. Only honored on the first message of the stream. + Append bool `protobuf:"varint,4,opt,name=append,proto3" json:"append,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -767,6 +770,13 @@ func (x *WriteFileRequest) GetMode() uint32 { return 0 } +func (x *WriteFileRequest) GetAppend() bool { + if x != nil { + return x.Append + } + return false +} + // Response confirming the write operation. type WriteFileResponse struct { state protoimpl.MessageState `protogen:"open.v1"` @@ -857,11 +867,12 @@ const file_proto_ateenv_v1alpha_guest_proto_rawDesc = "" + "\x0fReadFileRequest\x12\x12\n" + "\x04path\x18\x01 \x01(\tR\x04path\"\x1f\n" + "\tFileChunk\x12\x12\n" + - "\x04data\x18\x01 \x01(\fR\x04data\"P\n" + + "\x04data\x18\x01 \x01(\fR\x04data\"h\n" + "\x10WriteFileRequest\x12\x12\n" + "\x04path\x18\x01 \x01(\tR\x04path\x12\x14\n" + "\x05chunk\x18\x02 \x01(\fR\x05chunk\x12\x12\n" + - "\x04mode\x18\x03 \x01(\rR\x04mode\"8\n" + + "\x04mode\x18\x03 \x01(\rR\x04mode\x12\x16\n" + + "\x06append\x18\x04 \x01(\bR\x06append\"8\n" + "\x11WriteFileResponse\x12#\n" + "\rbytes_written\x18\x01 \x01(\x03R\fbytesWritten*\xa3\x01\n" + "\rProcessStatus\x12\x1e\n" + diff --git a/proto/ateenv/v1alpha/guest.proto b/proto/ateenv/v1alpha/guest.proto index b979b7a..1b6af4f 100644 --- a/proto/ateenv/v1alpha/guest.proto +++ b/proto/ateenv/v1alpha/guest.proto @@ -174,6 +174,9 @@ message WriteFileRequest { bytes chunk = 2; // Optional POSIX file mode permission (e.g. 0644 or 0755; processed on first message). uint32 mode = 3; + // Append to the file instead of truncating it. The file is created if it + // does not exist. Only honored on the first message of the stream. + bool append = 4; } // Response confirming the write operation.