diff --git a/platform/api/resource_handlers.go b/platform/api/resource_handlers.go index 8b9d578..283b3c1 100644 --- a/platform/api/resource_handlers.go +++ b/platform/api/resource_handlers.go @@ -95,6 +95,7 @@ func (h *coreHandlers) register(mux *http.ServeMux) { mux.HandleFunc("/api/v1/server-instances/{id}/logs/events", h.serverLogEvents) mux.HandleFunc("/api/v1/server-instances/{id}/files/workspace", h.serverFilesWorkspace) mux.HandleFunc("/api/v1/server-instances/{id}/files/list", h.serverFilesList) + mux.HandleFunc("/api/v1/server-instances/{id}/files/browse", h.serverFilesBrowse) mux.HandleFunc("/api/v1/server-instances/{id}/files/refresh", h.serverFilesRefresh) mux.HandleFunc("/api/v1/server-instances/{id}/files/read-snapshot", h.serverFilesReadSnapshot) mux.HandleFunc("/api/v1/server-instances/{id}/files/read", h.serverFilesRead) @@ -1331,6 +1332,44 @@ func (h *coreHandlers) serverFilesList(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, dto.ServerFileListFromDomain(result)) } +// serverFilesBrowse godoc +// @Summary Browse one live server directory +// @Description Dispatches a fresh Run files.list job for the requested logical directory and briefly waits for the matching result without exposing host paths or stale cached listings. +// @Tags server-files +// @Accept json +// @Produce json +// @Param id path string true "Server instance ID" +// @Param body body dto.ServerFileListRequest true "Server file browse request" +// @Success 202 {object} dto.ServerFileListResponse +// @Failure 400 {object} dto.ErrorResponse +// @Failure 401 {object} dto.ErrorResponse +// @Failure 403 {object} dto.ErrorResponse +// @Failure 404 {object} dto.ErrorResponse +// @Failure 405 {object} dto.ErrorResponse +// @Router /api/v1/server-instances/{id}/files/browse [post] +func (h *coreHandlers) serverFilesBrowse(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + writeMethodNotAllowed(w, http.MethodPost) + return + } + request, err := decodeJSON[dto.ServerFileListRequest](r) + if err != nil { + writeDecodeError(w, err) + return + } + request.DirectoryKey, err = h.serverFileDirectoryKey(r, request.DirectoryKey) + if err != nil { + writeServiceError(w, err) + return + } + result, err := h.core.BrowseServerFilesForSession(r.Context(), bearerToken(r), request.ToDomain(r.PathValue("id"))) + if err != nil { + writeServiceError(w, err) + return + } + writeJSON(w, http.StatusAccepted, dto.ServerFileListFromDomain(result)) +} + func (h *coreHandlers) serverFilesRefresh(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { writeMethodNotAllowed(w, http.MethodPost) @@ -1702,7 +1741,7 @@ func (h *coreHandlers) runJobClaim(w http.ResponseWriter, r *http.Request) { writeDecodeError(w, err) return } - result, err := h.core.ClaimRunJob(request.ToDomain()) + result, err := h.core.ClaimRunJobWithWait(r.Context(), request.ToDomain()) if err != nil { writeServiceError(w, err) return diff --git a/platform/api/routes.md b/platform/api/routes.md index f5cf8d7..8b78dd2 100644 --- a/platform/api/routes.md +++ b/platform/api/routes.md @@ -17,7 +17,7 @@ Routes use JSON request and response bodies unless a route explicitly accepts fi | Server runtime distribution | n/a | `GET /api/v1/server-instances/{id}/runtime/actions`, `POST /api/v1/server-instances/{id}/run/generate`, `POST /api/v1/server-instances/{id}/run/download`, `POST /api/v1/server-instances/{id}/run/key/reset`, `POST /api/v1/server-instances/{id}/run/update`, `GET /api/v1/server-instances/{id}/run/update`, `POST /api/v1/server-instances/{id}/client-managers/generate`, `POST /api/v1/server-instances/{id}/client-managers/download`, `POST /api/v1/server-instances/{id}/client-managers/key/reset`, `GET /api/v1/server-instances/{id}/dependencies`, `POST /api/v1/server-instances/{id}/dependencies/check`, `POST /api/v1/server-instances/{id}/dependencies/install` | `ServerRuntimeActionsResponse`, `RunDistributionGenerateRequest`, `RunDistributionResponse`, `RunUpdateRequest`, `RunUpdateJobResponse`/`RunUpdateJobListResponse`, `ClientManagerBuildRequest`, `ClientManagerDistributionResponse`, `ClientManagerDownloadRequest`, `ComponentKeyResetRequest`, `ComponentKeyResponse`, `DependencyCatalogResponse`, `DependencyJobRequest` | | Metrics | `GET /api/v1/metrics/platform`, `GET /api/v1/metrics/server-instances` | n/a | `PlatformResourceUsageResponse`, `ServerMetricsResponse`, `ServerMetricsListResponse` | | File operations | `POST /api/v1/file-operations/dispatch` | n/a | `FileOperationDispatchRequest`, `FileOperationDispatchResponse` | -| Server file manager | `GET /api/v1/server-instances/{id}/files/workspace`, `GET /api/v1/server-instances/{id}/files/list`, `POST /api/v1/server-instances/{id}/files/refresh`, `POST /api/v1/server-instances/{id}/files/read`, `POST /api/v1/server-instances/{id}/files/write`, `POST /api/v1/server-instances/{id}/files/upload`, `POST /api/v1/server-instances/{id}/files/download` | `GET /api/v1/server-instances/{id}/files/read-snapshot` | `ServerFileWorkspaceResponse`, `ServerFileListResponse`, `DeclaredFileReadSnapshotResponse`, `ServerFileReadRequest`, `ServerFileWriteRequest`, `ServerFileUploadResponse`, `ServerFileDownloadRequest`, `ServerFileDownloadResponse` | +| Server file manager | `GET /api/v1/server-instances/{id}/files/workspace`, `POST /api/v1/server-instances/{id}/files/browse`, `GET /api/v1/server-instances/{id}/files/list`, `POST /api/v1/server-instances/{id}/files/refresh`, `POST /api/v1/server-instances/{id}/files/read`, `POST /api/v1/server-instances/{id}/files/write`, `POST /api/v1/server-instances/{id}/files/upload`, `POST /api/v1/server-instances/{id}/files/download` | `GET /api/v1/server-instances/{id}/files/read-snapshot` | `ServerFileWorkspaceResponse`, `ServerFileListResponse`, `DeclaredFileReadSnapshotResponse`, `ServerFileReadRequest`, `ServerFileWriteRequest`, `ServerFileUploadResponse`, `ServerFileDownloadRequest`, `ServerFileDownloadResponse` | | Plugin-owned data | n/a | `GET/PUT/DELETE /api/v1/server-instances/{id}/plugin-data/{collection}`, `POST .../plugin-data/{collection}/transaction` | `PluginDataPutRequest`, `PluginDataTransactionRequest`, `PluginDataRecordResponse`, `PluginDataListResponse` | | Server administrators | `GET /api/v1/server-instances/{id}/administrators/candidates`, `POST /api/v1/server-instances/{id}/administrators` | `DELETE /api/v1/server-instances/{id}/administrators/{userId}` | `ServerMemberRequest`, `ServerMemberResponse`, `ServerMemberListResponse`, `ServerInstanceResponse` | | Run endpoints | `GET /api/v1/run/endpoints`, `POST /api/v1/run/endpoints` | `GET /api/v1/run/endpoints/{id}` | `RunEndpointCreateRequest`, `RunEndpointResponse`, `RunEndpointListResponse` | @@ -176,7 +176,7 @@ Control is the highest-priority run-facing channel; artifact/file transfer press ## Implemented Run Job Actions -- `POST /api/v1/run/jobs/claim`: accept `RunJobClaimRequest`, validate the active Run session, sweep expired endpoint work, and durably claim one eligible queued/retrying job with a monotonic per-job attempt, hashed lease credential, ack deadline, and execution lease. +- `POST /api/v1/run/jobs/claim`: accept `RunJobClaimRequest`, validate the active Run session, optionally hold the request for a bounded `waitSeconds` window, sweep expired endpoint work, and durably claim one eligible queued/retrying job with a monotonic per-job attempt, hashed lease credential, ack deadline, and execution lease. - `POST /api/v1/run/jobs/ack`: accept `RunJobAckRequest`, fence endpoint/session generation/attempt/lease, reject late acknowledgements, and move the current attempt into running state. - `POST /api/v1/run/jobs/progress`: accept `RunJobProgressRequest`, reject stale sequences and expired/old attempts, persist bounded progress, and renew the current execution lease. - `POST /api/v1/run/jobs/result`: accept `RunJobResultRequest` and write an idempotent terminal result or durable retry-wait transition with capped exponential backoff. diff --git a/platform/domain/job_channel.go b/platform/domain/job_channel.go index d9b622f..8888460 100644 --- a/platform/domain/job_channel.go +++ b/platform/domain/job_channel.go @@ -37,6 +37,7 @@ type RunJobClaim struct { SessionToken string Capabilities []string Capacity RunCapacity + WaitSeconds int } type RunJobClaimResult struct { diff --git a/platform/dto/job_channel.go b/platform/dto/job_channel.go index b744703..46d40bd 100644 --- a/platform/dto/job_channel.go +++ b/platform/dto/job_channel.go @@ -35,6 +35,7 @@ type RunJobClaimRequest struct { SessionToken string `json:"sessionToken"` Capabilities []string `json:"capabilities"` Capacity RunCapacityResponse `json:"capacity"` + WaitSeconds int `json:"waitSeconds,omitempty"` } type RunJobClaimResponse struct { @@ -405,6 +406,7 @@ func (request RunJobClaimRequest) ToDomain() domain.RunJobClaim { SessionToken: request.SessionToken, Capabilities: domain.CopyStringSlice(request.Capabilities), Capacity: capacityToDomain(request.Capacity), + WaitSeconds: request.WaitSeconds, } } diff --git a/platform/service/artifact_body_store.go b/platform/service/artifact_body_store.go index cb7a81a..8bc9782 100644 --- a/platform/service/artifact_body_store.go +++ b/platform/service/artifact_body_store.go @@ -6,6 +6,7 @@ import ( "encoding/json" "errors" "fmt" + "io" "os" "path/filepath" "sort" @@ -22,6 +23,8 @@ type ArtifactBodyStore interface { LoadTransfers() ([]domain.ArtifactTransferSession, error) PutPayload(string, []byte) error GetPayload(string) ([]byte, error) + ReadPayloadRange(string, int64, int) ([]byte, error) + CommitTransferPayload(domain.ArtifactTransferSession) error } type MemoryArtifactBodyStore struct { @@ -73,6 +76,40 @@ func (store *MemoryArtifactBodyStore) GetPayload(artifactID string) ([]byte, err return domain.CopyBytes(payload), nil } +func (store *MemoryArtifactBodyStore) ReadPayloadRange(artifactID string, offset int64, length int) ([]byte, error) { + store.mu.Lock() + defer store.mu.Unlock() + payload, exists := store.payloads[artifactID] + if !exists { + return nil, repo.ErrNotFound + } + if offset < 0 || offset > int64(len(payload)) || length < 0 || int64(length) > int64(len(payload))-offset { + return nil, fmt.Errorf("artifact range is invalid") + } + return domain.CopyBytes(payload[int(offset) : int(offset)+length]), nil +} + +func (store *MemoryArtifactBodyStore) CommitTransferPayload(session domain.ArtifactTransferSession) error { + store.mu.Lock() + defer store.mu.Unlock() + payload := make([]byte, 0, int(session.SizeBytes)) + for index := 0; index < session.TotalChunks; index++ { + record, exists := session.ReceivedChunks[index] + if !exists { + return validationError("artifact transfer has missing chunks") + } + payload = append(payload, record.Payload...) + } + if int64(len(payload)) != session.SizeBytes { + return validationError("artifact transfer size does not match metadata") + } + if checksum := validator.BytesChecksum(payload); checksum != session.Checksum { + return validationError("artifact transfer checksum does not match metadata") + } + store.payloads[session.ArtifactID] = payload + return nil +} + type FileArtifactBodyStore struct { mu sync.Mutex rootDir string @@ -101,12 +138,14 @@ func (store *FileArtifactBodyStore) SaveTransfer(session domain.ArtifactTransfer } manifest := domain.CopyArtifactTransferSession(session) for index, record := range manifest.ReceivedChunks { - payload := domain.CopyBytes(record.Payload) - if len(payload) != record.SizeBytes || validator.BytesChecksum(payload) != record.Checksum { - return validationError("artifact chunk does not match durable manifest") - } - if err := writeAtomicFile(filepath.Join(dir, fmt.Sprintf("chunk-%08d.bin", index)), payload, 0o600); err != nil { - return err + if record.Payload != nil { + payload := domain.CopyBytes(record.Payload) + if len(payload) != record.SizeBytes || validator.BytesChecksum(payload) != record.Checksum { + return validationError("artifact chunk does not match durable manifest") + } + if err := writeAtomicFile(filepath.Join(dir, fmt.Sprintf("chunk-%08d.bin", index)), payload, 0o600); err != nil { + return err + } } record.Payload = nil manifest.ReceivedChunks[index] = record @@ -148,14 +187,11 @@ func (store *FileArtifactBodyStore) LoadTransfers() ([]domain.ArtifactTransferSe return nil, fmt.Errorf("artifact transfer manifest identity mismatch") } for index, record := range session.ReceivedChunks { - payload, err := os.ReadFile(filepath.Join(dir, fmt.Sprintf("chunk-%08d.bin", index))) - if err != nil { - return nil, fmt.Errorf("read artifact transfer chunk: %w", err) + chunkPath := filepath.Join(dir, fmt.Sprintf("chunk-%08d.bin", index)) + if _, err := os.Stat(chunkPath); err != nil { + return nil, fmt.Errorf("stat artifact transfer chunk: %w", err) } - if len(payload) != record.SizeBytes || validator.BytesChecksum(payload) != record.Checksum { - return nil, validationError("durable artifact chunk checksum mismatch") - } - record.Payload = payload + record.Payload = nil session.ReceivedChunks[index] = record } out = append(out, domain.CopyArtifactTransferSession(session)) @@ -182,6 +218,107 @@ func (store *FileArtifactBodyStore) GetPayload(artifactID string) ([]byte, error return payload, nil } +func (store *FileArtifactBodyStore) ReadPayloadRange(artifactID string, offset int64, length int) ([]byte, error) { + if offset < 0 || length < 0 { + return nil, fmt.Errorf("artifact range is invalid") + } + store.mu.Lock() + defer store.mu.Unlock() + file, err := os.Open(store.payloadPath(artifactID)) + if errors.Is(err, os.ErrNotExist) { + return nil, repo.ErrNotFound + } + if err != nil { + return nil, fmt.Errorf("open artifact payload: %w", err) + } + defer file.Close() + payload := make([]byte, length) + read, err := file.ReadAt(payload, offset) + if err != nil && !(errors.Is(err, io.ErrUnexpectedEOF) && read == length) { + return nil, fmt.Errorf("read artifact payload range: %w", err) + } + if read != length { + return nil, fmt.Errorf("artifact payload range is shorter than requested") + } + return payload, nil +} + +func (store *FileArtifactBodyStore) CommitTransferPayload(session domain.ArtifactTransferSession) error { + store.mu.Lock() + defer store.mu.Unlock() + + transferDir := store.transferDir(session.TransferID) + payloadPath := store.payloadPath(session.ArtifactID) + if err := os.MkdirAll(filepath.Dir(payloadPath), 0o700); err != nil { + return fmt.Errorf("create artifact payload directory: %w", err) + } + tmp := payloadPath + ".tmp" + out, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600) + if err != nil { + return fmt.Errorf("open artifact payload temporary file: %w", err) + } + hash := sha256.New() + written := int64(0) + for index := 0; index < session.TotalChunks; index++ { + record, exists := session.ReceivedChunks[index] + if !exists { + _ = out.Close() + _ = os.Remove(tmp) + return validationError("artifact transfer has missing chunks") + } + chunkPath := filepath.Join(transferDir, fmt.Sprintf("chunk-%08d.bin", index)) + chunk, err := os.Open(chunkPath) + if err != nil { + _ = out.Close() + _ = os.Remove(tmp) + return fmt.Errorf("open artifact transfer chunk: %w", err) + } + chunkHash := sha256.New() + count, copyErr := io.Copy(io.MultiWriter(out, hash, chunkHash), chunk) + closeErr := chunk.Close() + if copyErr != nil { + _ = out.Close() + _ = os.Remove(tmp) + return fmt.Errorf("copy artifact transfer chunk: %w", copyErr) + } + if closeErr != nil { + _ = out.Close() + _ = os.Remove(tmp) + return fmt.Errorf("close artifact transfer chunk: %w", closeErr) + } + if int(count) != record.SizeBytes || "sha256:"+hex.EncodeToString(chunkHash.Sum(nil)) != record.Checksum { + _ = out.Close() + _ = os.Remove(tmp) + return validationError("durable artifact chunk checksum mismatch") + } + written += count + } + if written != session.SizeBytes { + _ = out.Close() + _ = os.Remove(tmp) + return validationError("artifact transfer size does not match metadata") + } + if checksum := "sha256:" + hex.EncodeToString(hash.Sum(nil)); checksum != session.Checksum { + _ = out.Close() + _ = os.Remove(tmp) + return validationError("artifact transfer checksum does not match metadata") + } + if err := out.Sync(); err != nil { + _ = out.Close() + _ = os.Remove(tmp) + return fmt.Errorf("sync artifact payload: %w", err) + } + if err := out.Close(); err != nil { + _ = os.Remove(tmp) + return fmt.Errorf("close artifact payload: %w", err) + } + if err := os.Rename(tmp, payloadPath); err != nil { + _ = os.Remove(tmp) + return fmt.Errorf("commit artifact payload: %w", err) + } + return nil +} + func (store *FileArtifactBodyStore) transferDir(transferID string) string { return filepath.Join(store.rootDir, "transfers", stableStorageKey(transferID)) } diff --git a/platform/service/artifact_download.go b/platform/service/artifact_download.go index 72d6471..e4d0044 100644 --- a/platform/service/artifact_download.go +++ b/platform/service/artifact_download.go @@ -14,6 +14,7 @@ import ( ) const artifactDownloadStorageBehavior = "platform-durable-artifact-store" +const artifactPayloadCacheLimit = 4 * 1024 * 1024 func (svc *CoreService) GetArtifactForSession(sessionID string, artifactID string) (domain.Artifact, error) { artifact, err := svc.store.Artifacts().Get(strings.TrimSpace(artifactID)) @@ -73,29 +74,22 @@ func (svc *CoreService) ReadArtifactContentForSession(sessionID string, request if artifact.State != domain.ArtifactStateAvailable { return domain.ArtifactContent{}, validationError("artifact must be available before download") } - payload, err := svc.artifactPayload(artifact.ID) - if err != nil { - return domain.ArtifactContent{}, err - } - if int64(len(payload)) != artifact.SizeBytes { - return domain.ArtifactContent{}, validationError("artifact content size does not match metadata") - } - if checksum := validator.BytesChecksum(payload); checksum != artifact.Checksum { - return domain.ArtifactContent{}, validationError("artifact content checksum does not match metadata") - } - if request.Offset >= artifact.SizeBytes { - return domain.ArtifactContent{}, validationError("offset must be inside artifact content") - } limit := request.Limit if limit == 0 { limit = validator.MaxArtifactDownloadBytes } + if request.Offset > artifact.SizeBytes { + return domain.ArtifactContent{}, validationError("artifact range exceeds metadata") + } remaining := artifact.SizeBytes - request.Offset if int64(limit) > remaining { limit = int(remaining) } - end := int(request.Offset) + limit - part := domain.CopyBytes(payload[int(request.Offset):end]) + payload, err := svc.artifactStore.ReadPayloadRange(artifact.ID, request.Offset, limit) + if err != nil { + return domain.ArtifactContent{}, err + } + part := domain.CopyBytes(payload) filename, contentType := svc.artifactDownloadPresentation(artifact) content := domain.ArtifactContent{ ArtifactID: artifact.ID, @@ -162,7 +156,7 @@ func (svc *CoreService) artifactPayload(artifactID string) ([]byte, error) { return domain.CopyBytes(payload), nil } if payload, err := svc.artifactStore.GetPayload(artifactID); err == nil { - svc.artifactPayloads[artifactID] = domain.CopyBytes(payload) + svc.cacheArtifactPayload(artifactID, payload) return payload, nil } else if !errors.Is(err, repo.ErrNotFound) { return nil, err @@ -193,10 +187,18 @@ func (svc *CoreService) artifactPayload(artifactID string) ([]byte, error) { if err := svc.artifactStore.PutPayload(artifactID, payload); err != nil { return nil, err } - svc.artifactPayloads[artifactID] = domain.CopyBytes(payload) + svc.cacheArtifactPayload(artifactID, payload) return payload, nil } +func (svc *CoreService) cacheArtifactPayload(artifactID string, payload []byte) { + if len(payload) > artifactPayloadCacheLimit { + delete(svc.artifactPayloads, artifactID) + return + } + svc.artifactPayloads[artifactID] = domain.CopyBytes(payload) +} + func artifactDownloadFilename(artifactID string) string { name := strings.TrimSpace(artifactID) if name == "" || strings.Contains(name, "/") || strings.Contains(name, `\`) || strings.Contains(name, "://") { diff --git a/platform/service/artifact_transfer.go b/platform/service/artifact_transfer.go index 163ba83..fd05a04 100644 --- a/platform/service/artifact_transfer.go +++ b/platform/service/artifact_transfer.go @@ -1,7 +1,6 @@ package service import ( - "bytes" "errors" "fmt" "sort" @@ -111,7 +110,7 @@ func (svc *CoreService) UploadArtifactChunk(chunk domain.ArtifactChunkUpload) (d } if existing, exists := session.ReceivedChunks[chunk.ChunkIndex]; exists { - if existing.Offset == chunk.Offset && existing.SizeBytes == chunk.SizeBytes && existing.Checksum == chunk.Checksum && bytes.Equal(existing.Payload, chunk.Payload) { + if existing.Offset == chunk.Offset && existing.SizeBytes == chunk.SizeBytes && existing.Checksum == chunk.Checksum { return artifactChunkUploadResult(session, chunk.ChunkIndex, true, stamp), nil } return domain.ArtifactChunkUploadResult{}, validationError("artifact chunk conflicts with acknowledged chunk") @@ -129,7 +128,7 @@ func (svc *CoreService) UploadArtifactChunk(chunk domain.ArtifactChunkUpload) (d if err := svc.artifactStore.SaveTransfer(session); err != nil { return domain.ArtifactChunkUploadResult{}, err } - svc.artifactTransfers[session.TransferID] = domain.CopyArtifactTransferSession(session) + svc.artifactTransfers[session.TransferID] = svc.artifactTransferSessionForMemory(session) return artifactChunkUploadResult(session, chunk.ChunkIndex, false, stamp), nil } @@ -185,19 +184,13 @@ func (svc *CoreService) CompleteArtifactTransfer(complete domain.ArtifactTransfe return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer has missing chunks") } - payload := make([]byte, 0, int(session.SizeBytes)) for index := 0; index < session.TotalChunks; index++ { - record, exists := session.ReceivedChunks[index] - if !exists { + if _, exists := session.ReceivedChunks[index]; !exists { return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer has missing chunks") } - payload = append(payload, record.Payload...) } - if int64(len(payload)) != session.SizeBytes { - return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer size does not match metadata") - } - if checksum := validator.BytesChecksum(payload); checksum != session.Checksum { - return domain.ArtifactTransferCompleteResult{}, validationError("artifact transfer checksum does not match metadata") + if err := svc.artifactStore.CommitTransferPayload(session); err != nil { + return domain.ArtifactTransferCompleteResult{}, err } artifact.SizeBytes = session.SizeBytes @@ -207,22 +200,29 @@ func (svc *CoreService) CompleteArtifactTransfer(complete domain.ArtifactTransfe if err := validator.ValidateArtifact(artifact); err != nil { return domain.ArtifactTransferCompleteResult{}, err } - if err := svc.artifactStore.PutPayload(artifact.ID, payload); err != nil { - return domain.ArtifactTransferCompleteResult{}, err - } if err := svc.store.Artifacts().Update(artifact); err != nil { return domain.ArtifactTransferCompleteResult{}, err } - svc.artifactPayloads[artifact.ID] = domain.CopyBytes(payload) session.Completed = true session.UpdatedAt = stamp if err := svc.artifactStore.SaveTransfer(session); err != nil { return domain.ArtifactTransferCompleteResult{}, err } - svc.artifactTransfers[session.TransferID] = domain.CopyArtifactTransferSession(session) + svc.artifactTransfers[session.TransferID] = svc.artifactTransferSessionForMemory(session) return domain.ArtifactTransferCompleteResult{Accepted: true, TransferID: session.TransferID, Artifact: artifact, Completed: true, ServerTime: stamp}, nil } +func (svc *CoreService) artifactTransferSessionForMemory(session domain.ArtifactTransferSession) domain.ArtifactTransferSession { + cached := domain.CopyArtifactTransferSession(session) + if _, durableFileStore := svc.artifactStore.(*FileArtifactBodyStore); durableFileStore { + for index, record := range cached.ReceivedChunks { + record.Payload = nil + cached.ReceivedChunks[index] = record + } + } + return cached +} + func (svc *CoreService) getArtifactTransferSession(transferID string) (domain.ArtifactTransferSession, error) { session, exists := svc.artifactTransfers[transferID] if !exists { diff --git a/platform/service/distribution_build_execution.go b/platform/service/distribution_build_execution.go index 6556ed8..13d37e9 100644 --- a/platform/service/distribution_build_execution.go +++ b/platform/service/distribution_build_execution.go @@ -433,7 +433,7 @@ func (svc *CoreService) storeDistributionBuildArtifact(artifactID string, jobID } else if err := svc.store.Artifacts().Update(artifact); err != nil { return err } - svc.artifactPayloads[artifact.ID] = domain.CopyBytes(payload) + svc.cacheArtifactPayload(artifact.ID, payload) return nil } diff --git a/platform/service/distributions.go b/platform/service/distributions.go index 87e0a2d..7196aa0 100644 --- a/platform/service/distributions.go +++ b/platform/service/distributions.go @@ -932,7 +932,7 @@ func (svc *CoreService) ensureArtifactPayload(artifactID string, payload []byte, if int64(len(existingPayload)) != artifact.SizeBytes || validator.BytesChecksum(existingPayload) != artifact.Checksum { return validationError("artifact payload does not match metadata") } - svc.artifactPayloads[artifactID] = domain.CopyBytes(existingPayload) + svc.cacheArtifactPayload(artifactID, existingPayload) return nil } else if !errors.Is(err, repo.ErrNotFound) { return err @@ -943,7 +943,7 @@ func (svc *CoreService) ensureArtifactPayload(artifactID string, payload []byte, if err := svc.artifactStore.PutPayload(artifactID, payload); err != nil { return err } - svc.artifactPayloads[artifactID] = domain.CopyBytes(payload) + svc.cacheArtifactPayload(artifactID, payload) return nil } diff --git a/platform/service/durable_observability_test.go b/platform/service/durable_observability_test.go index c1b5969..8ffa8de 100644 --- a/platform/service/durable_observability_test.go +++ b/platform/service/durable_observability_test.go @@ -1,7 +1,6 @@ package service import ( - "bytes" "path/filepath" "testing" "time" @@ -91,8 +90,8 @@ func TestFileArtifactBodyStoreResumesTransferAfterServiceRestart(t *testing.T) { if err != nil { t.Fatalf("final restart service: %v", err) } - stored, err := finalService.artifactPayload("durable-artifact") - if err != nil || !bytes.Equal(stored, payload) { + stored, err := finalService.artifactStore.ReadPayloadRange("durable-artifact", 0, len(payload)) + if err != nil || string(stored) != string(payload) { t.Fatalf("expected durable payload after restart, payload=%q err=%v", stored, err) } } diff --git a/platform/service/job_channel.go b/platform/service/job_channel.go index 67262dc..ac56713 100644 --- a/platform/service/job_channel.go +++ b/platform/service/job_channel.go @@ -1,6 +1,7 @@ package service import ( + "context" "crypto/subtle" "fmt" "sort" @@ -85,6 +86,68 @@ func (svc *CoreService) ClaimRunJob(claim domain.RunJobClaim) (domain.RunJobClai }), nil } +func (svc *CoreService) ClaimRunJobWithWait(ctx context.Context, claim domain.RunJobClaim) (domain.RunJobClaimResult, error) { + claim = domain.CopyRunJobClaim(claim) + result, err := svc.ClaimRunJob(claim) + if err != nil || result.HasJob || claim.WaitSeconds <= 0 || runJobClaimAtCapacity(claim) { + return result, err + } + waiter := svc.registerRunJobWaiter(claim.RunEndpointID) + defer svc.unregisterRunJobWaiter(claim.RunEndpointID, waiter) + result, err = svc.ClaimRunJob(claim) + if err != nil || result.HasJob { + return result, err + } + timer := time.NewTimer(time.Duration(claim.WaitSeconds) * time.Second) + defer timer.Stop() + select { + case <-ctx.Done(): + return domain.RunJobClaimResult{}, ctx.Err() + case <-waiter: + case <-timer.C: + } + return svc.ClaimRunJob(claim) +} + +func runJobClaimAtCapacity(claim domain.RunJobClaim) bool { + return claim.Capacity.MaxJobs > 0 && claim.Capacity.RunningJobs >= claim.Capacity.MaxJobs +} + +func (svc *CoreService) registerRunJobWaiter(runEndpointID string) chan struct{} { + waiter := make(chan struct{}) + svc.jobWaitMu.Lock() + svc.jobWaiters[runEndpointID] = append(svc.jobWaiters[runEndpointID], waiter) + svc.jobWaitMu.Unlock() + return waiter +} + +func (svc *CoreService) unregisterRunJobWaiter(runEndpointID string, waiter chan struct{}) { + svc.jobWaitMu.Lock() + waiters := svc.jobWaiters[runEndpointID] + for index, candidate := range waiters { + if candidate == waiter { + waiters = append(waiters[:index], waiters[index+1:]...) + break + } + } + if len(waiters) == 0 { + delete(svc.jobWaiters, runEndpointID) + } else { + svc.jobWaiters[runEndpointID] = waiters + } + svc.jobWaitMu.Unlock() +} + +func (svc *CoreService) notifyRunJobWaiters(runEndpointID string) { + svc.jobWaitMu.Lock() + waiters := svc.jobWaiters[runEndpointID] + delete(svc.jobWaiters, runEndpointID) + svc.jobWaitMu.Unlock() + for _, waiter := range waiters { + close(waiter) + } +} + func withoutCapability(capabilities []string, forbidden string) []string { filtered := make([]string, 0, len(capabilities)) for _, capability := range capabilities { @@ -570,7 +633,11 @@ func (svc *CoreService) updateScheduledJob(job domain.Job) error { if err := validator.ValidateJob(job); err != nil { return err } - return svc.store.Jobs().Update(job) + if err := svc.store.Jobs().Update(job); err != nil { + return err + } + svc.notifyRunJobWaiters(job.RunEndpointID) + return nil } func normalizeJobScheduling(job domain.Job, stamp time.Time) domain.Job { diff --git a/platform/service/job_channel_test.go b/platform/service/job_channel_test.go index 0941279..27a6d81 100644 --- a/platform/service/job_channel_test.go +++ b/platform/service/job_channel_test.go @@ -1,8 +1,10 @@ package service import ( + "context" "strings" "testing" + "time" "browser.local/platform/domain" ) @@ -98,6 +100,34 @@ func TestCoreServiceRunJobClaimNoJob(t *testing.T) { } } +func TestCoreServiceRunJobClaimWithWaitWakesOnCreate(t *testing.T) { + svc, sessionToken := newRegisteredRunJobService(t) + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + resultCh := make(chan domain.RunJobClaimResult, 1) + errCh := make(chan error, 1) + go func() { + claim, err := svc.ClaimRunJobWithWait(ctx, domain.RunJobClaim{RunEndpointID: "run-local", SessionToken: sessionToken, Capabilities: []string{"process.start"}, Capacity: domain.RunCapacity{MaxJobs: 4}, WaitSeconds: 10}) + if err != nil { + errCh <- err + return + } + resultCh <- claim + }() + waitForRunJobWaiter(t, svc, "run-local") + createQueuedRunJob(t, svc, "job-wait", "idem-wait") + select { + case err := <-errCh: + t.Fatalf("claim wait failed: %v", err) + case claim := <-resultCh: + if !claim.HasJob || claim.Job == nil || claim.Job.JobID != "job-wait" { + t.Fatalf("expected wait claim to return created job, got %+v", claim) + } + case <-ctx.Done(): + t.Fatalf("claim wait timed out: %v", ctx.Err()) + } +} + func TestCoreServiceRunJobClaimSkipsServerFileCapabilityWithoutDeclaration(t *testing.T) { svc := newTestCoreService() plugin, endpoint := createPluginAndRunEndpoint(t, svc) @@ -314,3 +344,18 @@ func createQueuedRunJob(t *testing.T, svc *CoreService, id string, idempotencyKe } return job } + +func waitForRunJobWaiter(t *testing.T, svc *CoreService, runEndpointID string) { + t.Helper() + deadline := time.Now().Add(500 * time.Millisecond) + for time.Now().Before(deadline) { + svc.jobWaitMu.Lock() + count := len(svc.jobWaiters[runEndpointID]) + svc.jobWaitMu.Unlock() + if count > 0 { + return + } + time.Sleep(time.Millisecond) + } + t.Fatalf("waiter for %s was not registered", runEndpointID) +} diff --git a/platform/service/resources.go b/platform/service/resources.go index 32b628e..e014100 100644 --- a/platform/service/resources.go +++ b/platform/service/resources.go @@ -1,6 +1,7 @@ package service import ( + "context" "crypto/pbkdf2" "crypto/rand" "crypto/sha256" @@ -126,6 +127,7 @@ type Core interface { GetServerFileWorkspaceForSession(string, string) (domain.ServerFileWorkspaceView, error) ListServerFilesForSession(string, domain.ServerFileListRequest) (domain.ServerFileListResult, error) RefreshServerFileListForSession(string, domain.ServerFileListRequest) (domain.ServerFileListResult, error) + BrowseServerFilesForSession(context.Context, string, domain.ServerFileListRequest) (domain.ServerFileListResult, error) ReadServerFileForSession(string, domain.ServerFileReadRequest) (domain.FileOperationDispatchResult, error) WriteServerFileForSession(string, domain.ServerFileWriteRequest) (domain.FileOperationDispatchResult, error) UploadServerFileForSession(string, domain.ServerFileUploadRequest) (domain.ServerFileUploadDispatch, error) @@ -140,6 +142,7 @@ type Core interface { ListJobsForSession(string, domain.JobFilter) ([]domain.Job, error) RequestRunJobCancelForSession(string, domain.RunJobCancelRequest) (domain.RunJobCancelRequestResult, error) ClaimRunJob(domain.RunJobClaim) (domain.RunJobClaimResult, error) + ClaimRunJobWithWait(context.Context, domain.RunJobClaim) (domain.RunJobClaimResult, error) AckRunJob(domain.RunJobAck) (domain.RunJobAckResult, error) UpdateRunJobProgress(domain.RunJobProgress) (domain.RunJobProgressResult, error) CompleteRunJob(domain.RunJobResult) (domain.RunJobResultResult, error) @@ -231,6 +234,8 @@ type CoreService struct { runSessions map[string]domain.RunControlSession runSessionSeq uint64 jobMu sync.Mutex + jobWaitMu sync.Mutex + jobWaiters map[string][]chan struct{} bridgeMu sync.Mutex bridgeSeq uint64 logStore LogBodyStore @@ -279,6 +284,7 @@ func newCoreServiceWithLogStore(store repo.Store, logStore LogBodyStore, now fun now: now, authSessions: map[string]string{}, runSessions: map[string]domain.RunControlSession{}, + jobWaiters: map[string][]chan struct{}{}, logStore: logStore, logProjectionStates: map[string]map[string]pluginLogSequenceState{}, logEventSubscribers: map[uint64]logEventSubscriber{}, @@ -2516,6 +2522,9 @@ func (svc *CoreService) CreateJob(job domain.Job) (domain.Job, error) { if err := svc.ensureJobLogStreams(existing, stamp); err != nil { return domain.Job{}, err } + if !isTerminalJobState(existing.State) { + svc.notifyRunJobWaiters(existing.RunEndpointID) + } return existing, nil } if !errors.Is(err, repo.ErrNotFound) { @@ -2555,6 +2564,7 @@ func (svc *CoreService) CreateJob(job domain.Job) (domain.Job, error) { if err := svc.ensureJobLogStreams(job, stamp); err != nil { return domain.Job{}, err } + svc.notifyRunJobWaiters(job.RunEndpointID) return domain.CopyJob(job), nil } diff --git a/platform/service/resources_test.go b/platform/service/resources_test.go index 80e15d5..8e1a9f7 100644 --- a/platform/service/resources_test.go +++ b/platform/service/resources_test.go @@ -1,8 +1,10 @@ package service import ( + "context" "encoding/base64" "errors" + "fmt" "strings" "testing" "time" @@ -1004,6 +1006,70 @@ func TestServerFileListReportsFailedRuntimeRefresh(t *testing.T) { } } +func TestServerFileBrowseWaitsForFreshRunResultWithoutCachedList(t *testing.T) { + svc := newTestCoreService() + plugin, endpoint := createPluginAndRunEndpoint(t, svc) + endpoint.Capabilities = append(endpoint.Capabilities, domain.JobCapabilityFilesList) + if err := svc.store.RunEndpoints().Update(endpoint); err != nil { + t.Fatalf("update file list capability: %v", err) + } + ownerSession := createServiceUserAndLogin(t, svc, domain.User{ID: "user-file-browse", DisplayName: "File Browse", Email: "file-browse@example.test", Roles: []string{"server-owner"}, PasswordHash: "secret-password"}) + instance, err := svc.CreateServerInstanceForSession(ownerSession, domain.ServerInstance{ID: "server-file-browse", PluginID: plugin.ID, RunEndpointID: endpoint.ID, Name: "File Browse Server", State: domain.ServerInstanceStateRunning, Deployment: domain.ServerDeploymentDefinition{Mode: domain.ServerDeploymentModeGuided, ServerRoot: `C:\scumserver`}}) + if err != nil { + t.Fatalf("create server: %v", err) + } + createCompleteRuntimeBinding(t, svc, instance, "local") + oldJob, err := svc.CreateJob(domain.Job{ID: "job-file-list-old", ServerInstanceID: instance.ID, RunEndpointID: endpoint.ID, Capability: domain.JobCapabilityFilesList, TargetKey: "server-root", IdempotencyKey: "idem-file-list-old"}) + if err != nil { + t.Fatalf("create old list job: %v", err) + } + oldJob.State = domain.JobStateSucceeded + oldJob.ExecutionResult = domain.JobExecutionResult{Kind: "file.list", Content: runFileListFixture("server-root", "", ".platform")} + oldJob.TerminalAt = fixedTime.Add(10 * time.Minute) + oldJob.UpdatedAt = oldJob.TerminalAt + if err := svc.store.Jobs().Update(oldJob); err != nil { + t.Fatalf("store old list result: %v", err) + } + helloRequest := validRunControlHello() + helloRequest.CapabilityReport.Capabilities = append(helloRequest.CapabilityReport.Capabilities, domain.JobCapabilityFilesList) + helloRequest.CapabilityReport.Fingerprint = "cap-file-browse" + hello, err := svc.RegisterRunHello(helloRequest) + if err != nil { + t.Fatalf("register run hello: %v", err) + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + resultCh := make(chan domain.ServerFileListResult, 1) + errCh := make(chan error, 1) + go func() { + result, err := svc.BrowseServerFilesForSession(ctx, ownerSession, domain.ServerFileListRequest{ServerInstanceID: instance.ID, DirectoryKey: "server-root", IdempotencyKey: "idem-file-browse-fresh"}) + if err != nil { + errCh <- err + return + } + resultCh <- result + }() + job := waitForServerFileListJob(t, svc, instance.ID, "job-file-list-old") + claim, err := svc.ClaimRunJob(domain.RunJobClaim{RunEndpointID: endpoint.ID, SessionToken: hello.SessionToken, Capabilities: []string{domain.JobCapabilityFilesList}, Capacity: domain.RunCapacity{MaxJobs: 4}}) + if err != nil || !claim.HasJob || claim.Job.JobID != job.ID { + t.Fatalf("claim fresh file list job: claim=%+v err=%v", claim, err) + } + _, err = svc.CompleteRunJob(domain.RunJobResult{RunEndpointID: endpoint.ID, SessionToken: hello.SessionToken, JobID: claim.Job.JobID, LeaseToken: claim.Job.LeaseToken, Attempt: claim.Job.Attempt, State: domain.JobStateSucceeded, Progress: domain.RunJobProgressReport{Percent: 100, Message: "listed"}, ExecutionResult: domain.JobExecutionResult{Kind: "file.list", Content: runFileListFixture("server-root", "", "SCUM")}}) + if err != nil { + t.Fatalf("complete fresh file list job: %v", err) + } + select { + case err := <-errCh: + t.Fatalf("browse failed: %v", err) + case result := <-resultCh: + if result.State != "ready" || len(result.Entries) != 1 || result.Entries[0].Name != "SCUM" || result.Job.ID != job.ID { + t.Fatalf("expected fresh browse result, got %+v", result) + } + case <-ctx.Done(): + t.Fatalf("browse timed out: %v", ctx.Err()) + } +} + func TestServerFileListFallsBackToPluginWorkspaceWithoutRunListCapability(t *testing.T) { svc := newTestCoreService() plugin, endpoint := createPluginAndRunEndpoint(t, svc) @@ -2062,6 +2128,29 @@ func createCompleteRuntimeBinding(t *testing.T, svc *CoreService, instance domai return binding } +func waitForServerFileListJob(t *testing.T, svc *CoreService, serverInstanceID string, excludedJobID string) domain.Job { + t.Helper() + deadline := time.Now().Add(500 * time.Millisecond) + for time.Now().Before(deadline) { + jobs, err := svc.store.Jobs().List(domain.JobFilter{ServerInstanceID: serverInstanceID}) + if err != nil { + t.Fatalf("list server jobs: %v", err) + } + for _, job := range jobs { + if job.ID != excludedJobID && job.Capability == domain.JobCapabilityFilesList { + return job + } + } + time.Sleep(time.Millisecond) + } + t.Fatalf("fresh file list job was not created") + return domain.Job{} +} + +func runFileListFixture(directoryKey string, relativePath string, name string) string { + return fmt.Sprintf(`{"directoryKey":%q,"path":%q,"entries":[{"name":%q,"kind":"directory","directoryKey":%q,"relativePath":%q}]}`, directoryKey, relativePath, name, directoryKey, name) +} + func runtimeBindingTestKeyIsSensitive(key string) bool { normalized := strings.ToLower(key) return strings.Contains(normalized, "password") || strings.Contains(normalized, "credential") || strings.Contains(normalized, "secret") || strings.Contains(normalized, "token") || strings.Contains(normalized, "dsn") diff --git a/platform/service/server_files.go b/platform/service/server_files.go index f549887..3ce53ef 100644 --- a/platform/service/server_files.go +++ b/platform/service/server_files.go @@ -1,6 +1,7 @@ package service import ( + "context" "encoding/json" "errors" "fmt" @@ -19,6 +20,7 @@ const ( serverFileTransferChannel = "run-file-transfer" serverFileMaxInlineEditBytes = 64 * 1024 serverFileDefaultDirectoryKey = "server-root" + serverFileBrowseWait = 4 * time.Second ) type serverFileContext struct { @@ -188,6 +190,82 @@ func (svc *CoreService) RefreshServerFileListForSession(sessionID string, reques return domain.CopyServerFileListResult(domain.ServerFileListResult{ServerInstanceID: ctx.Instance.ID, PluginID: ctx.Plugin.ID, DirectoryKey: request.DirectoryKey, Path: request.Path, State: "pending", Entries: entries, Job: job, Reason: "目录刷新任务已派发到 Run。"}), nil } +func (svc *CoreService) BrowseServerFilesForSession(ctx context.Context, sessionID string, request domain.ServerFileListRequest) (domain.ServerFileListResult, error) { + request = normalizeServerFileListRequest(request) + if request.IdempotencyKey == "" { + idempotencyKey, err := svc.serverFileBrowseIdempotencyKey(request) + if err != nil { + return domain.ServerFileListResult{}, err + } + request.IdempotencyKey = idempotencyKey + } + result, err := svc.RefreshServerFileListForSession(sessionID, request) + if err != nil || result.State != "pending" || result.Job.ID == "" { + return result, err + } + timer := time.NewTimer(serverFileBrowseWait) + defer timer.Stop() + for { + current, ready, err := svc.serverFileListResultFromJob(request, result, result.Job.ID) + if err != nil || ready { + return current, err + } + waiter := svc.registerRunJobWaiter(result.Job.RunEndpointID) + current, ready, err = svc.serverFileListResultFromJob(request, result, result.Job.ID) + if err != nil || ready { + svc.unregisterRunJobWaiter(result.Job.RunEndpointID, waiter) + return current, err + } + select { + case <-ctx.Done(): + svc.unregisterRunJobWaiter(result.Job.RunEndpointID, waiter) + return domain.ServerFileListResult{}, ctx.Err() + case <-waiter: + case <-timer.C: + svc.unregisterRunJobWaiter(result.Job.RunEndpointID, waiter) + current.Reason = "正在读取目录。" + return current, nil + } + } +} + +func (svc *CoreService) serverFileListResultFromJob(request domain.ServerFileListRequest, fallback domain.ServerFileListResult, jobID string) (domain.ServerFileListResult, bool, error) { + job, err := svc.store.Jobs().Get(jobID) + if err != nil { + return domain.ServerFileListResult{}, true, err + } + result := domain.CopyServerFileListResult(fallback) + result.Job = job + if !isTerminalJobState(job.State) { + return result, false, nil + } + result.RefreshedAt = job.TerminalAt + if job.State != domain.JobStateSucceeded || job.ExecutionResult.Kind != "file.list" { + result.State = "failed" + result.Reason = serverFileListJobFailureReason(job) + return result, true, nil + } + entries, err := serverFileEntriesFromRunList(job.ExecutionResult.Content, request.DirectoryKey, request.Path) + if err != nil { + result.State = "failed" + result.Reason = "Run 返回的文件列表无法解析。" + return result, true, nil + } + result.State = "ready" + result.Entries = filterServerFileEntries(entries, request.Query) + result.Reason = "目录读取完成。" + return result, true, nil +} + +func (svc *CoreService) serverFileBrowseIdempotencyKey(request domain.ServerFileListRequest) (string, error) { + token, err := randomToken() + if err != nil { + return "", err + } + raw := strings.Join([]string{request.ServerInstanceID, request.DirectoryKey, request.Path, request.Query, strconv.FormatBool(request.Recursive), token}, ":") + return fmt.Sprintf("file-browse:%d", stableStringNumber(raw)), nil +} + func (svc *CoreService) ReadServerFileForSession(sessionID string, request domain.ServerFileReadRequest) (domain.FileOperationDispatchResult, error) { if err := validator.ValidateServerFileReadRequest(request); err != nil { return domain.FileOperationDispatchResult{}, err @@ -545,7 +623,7 @@ func (svc *CoreService) putBrowserFileArtifact(artifact domain.Artifact, payload if err := svc.artifactStore.PutPayload(artifact.ID, payload); err != nil { return err } - svc.artifactPayloads[artifact.ID] = domain.CopyBytes(payload) + svc.cacheArtifactPayload(artifact.ID, payload) return nil } diff --git a/platform/validator/job_channel.go b/platform/validator/job_channel.go index 99de394..ccee2da 100644 --- a/platform/validator/job_channel.go +++ b/platform/validator/job_channel.go @@ -8,6 +8,7 @@ import ( ) const maxJobChannelMessageLength = 256 +const maxRunJobClaimWaitSeconds = 30 func ValidateRunJobClaim(claim domain.RunJobClaim) error { var violations []string @@ -19,6 +20,9 @@ func ValidateRunJobClaim(claim domain.RunJobClaim) error { violations = append(violations, fmt.Sprintf("capabilities[%d] is required", i)) } } + if claim.WaitSeconds < 0 || claim.WaitSeconds > maxRunJobClaimWaitSeconds { + violations = append(violations, fmt.Sprintf("waitSeconds must be between 0 and %d", maxRunJobClaimWaitSeconds)) + } return finish(violations) } diff --git a/platform_web/api/client.test.ts b/platform_web/api/client.test.ts index 92573ad..262db33 100644 --- a/platform_web/api/client.test.ts +++ b/platform_web/api/client.test.ts @@ -281,6 +281,10 @@ describe("PlatformApiClient AI providers", () => { expect(JSON.parse(String(init.body))).toEqual({ directoryKey: "configs", query: "server", recursive: true, idempotencyKey: "idem-file-list" }); return jsonResponse({ serverInstanceId: server.id, pluginId: plugin.id, directoryKey: "configs", state: "pending", entries: [], job: { ...job, id: "job-file-list", capability: "files.list", targetKey: "configs" }, reason: "queued" }); } + if (url.endsWith("/api/v1/server-instances/server-1/files/browse") && init?.method === "POST") { + expect(JSON.parse(String(init.body))).toEqual({ directoryKey: "configs", query: "server", recursive: true, idempotencyKey: "idem-file-browse" }); + return jsonResponse({ serverInstanceId: server.id, pluginId: plugin.id, directoryKey: "configs", state: "ready", entries: [{ name: "server.properties", kind: "file", directoryKey: "configs", relativePath: "config/server.properties", logicalKey: "config/server.properties", scope: "config", sizeBytes: 42, editable: true, downloadable: true, remark: "配置文件" }], job: { ...job, id: "job-file-browse", capability: "files.list", targetKey: "configs" }, reason: "目录读取完成。" }); + } if (url.endsWith("/api/v1/server-instances/server-1/files/read-snapshot?key=config%2Fserver.properties")) { return jsonResponse({ serverInstanceId: server.id, pluginId: plugin.id, key: "config/server.properties", state: "ready", content: "server.name=Example\n", version: 3, checksum: "sha256:filechecksum", sizeBytes: 20, readAt: "2026-07-03T00:00:00Z" }); } @@ -595,6 +599,7 @@ describe("PlatformApiClient AI providers", () => { await expect(client.getServerFileWorkspace(server.id)).resolves.toMatchObject({ defaultDirectoryKey: "configs", transfer: { channel: "run-file-transfer" } }); await expect(client.listServerFiles(server.id, { directoryKey: "configs", query: "server", recursive: true })).resolves.toMatchObject({ state: "declared", entries: [{ logicalKey: "config/server.properties" }] }); await expect(client.refreshServerFiles(server.id, { directoryKey: "configs", query: "server", recursive: true, idempotencyKey: "idem-file-list" })).resolves.toMatchObject({ state: "pending", job: { capability: "files.list" } }); + await expect(client.browseServerFiles(server.id, { directoryKey: "configs", query: "server", recursive: true, idempotencyKey: "idem-file-browse" })).resolves.toMatchObject({ state: "ready", entries: [{ logicalKey: "config/server.properties" }] }); await expect(client.getServerFileReadSnapshot(server.id, "config/server.properties")).resolves.toMatchObject({ state: "ready", content: "server.name=Example\n" }); await expect(client.readServerFile(server.id, { key: "config/server.properties", idempotencyKey: "idem-file-read" })).resolves.toMatchObject({ operation: "read", job: { capability: "files.read" } }); await expect(client.writeServerFile(server.id, { key: "config/server.properties", content: "server.name=Example\n", expectedVersion: 3, expectedChecksum: "sha256:filechecksum", idempotencyKey: "idem-file-write" })).resolves.toMatchObject({ operation: "write", job: { capability: "files.write" } }); @@ -654,7 +659,7 @@ describe("PlatformApiClient AI providers", () => { client.invokeAI({ requestId: "ai-1", serverInstanceId: server.id, purpose: "config.suggest", prompt: "Tune PVP safely", currentConfig: "server.name=Example Survival #1\n" }) ).resolves.toMatchObject({ status: "ok", usage: { mocked: true }, configRecommendation: { diffSummary: "review required" } }); - expect(fetchMock).toHaveBeenCalledTimes(44); + expect(fetchMock).toHaveBeenCalledTimes(45); }); it("normalizes server file workspace null arrays from older platform responses", async () => { diff --git a/platform_web/api/client.ts b/platform_web/api/client.ts index f5415f1..c5128a5 100644 --- a/platform_web/api/client.ts +++ b/platform_web/api/client.ts @@ -596,6 +596,10 @@ export class PlatformApiClient { return normalizeServerFileList(await this.request(`/server-instances/${encodeURIComponent(serverInstanceId)}/files/refresh`, { method: "POST", body: request })); } + async browseServerFiles(serverInstanceId: string, request: ServerFileListRequest): Promise { + return normalizeServerFileList(await this.request(`/server-instances/${encodeURIComponent(serverInstanceId)}/files/browse`, { method: "POST", body: request })); + } + async readServerFile(serverInstanceId: string, request: ServerFileReadRequest): Promise { return this.request(`/server-instances/${encodeURIComponent(serverInstanceId)}/files/read`, { method: "POST", body: request }); } diff --git a/platform_web/api/contracts.md b/platform_web/api/contracts.md index 736eb69..51e9edd 100644 --- a/platform_web/api/contracts.md +++ b/platform_web/api/contracts.md @@ -28,7 +28,7 @@ Normal browser login uses the platform's HttpOnly SameSite cookie and `credentia - `listServerAdministratorCandidates`, `addServerAdministrator`, and `removeServerAdministrator` call server membership endpoints so server owners can invite or remove active non-platform-admin server administrators. - Game-specific pages use the scoped `plugin-data` collection API and declared plugin bridge machine actions; Platform does not expose game-specific projection or workflow clients. - `dispatchFileOperation` posts `FileOperationDispatchRequest` to `/file-operations/dispatch` using logical file keys and scoped refs rather than raw host paths; it remains the low-level compatibility dispatch for file work. -- `getServerFileWorkspace`, `listServerFiles`, `refreshServerFiles`, `readServerFile`, `getServerFileReadSnapshot`, `writeServerFile`, `uploadServerFile`, and `prepareServerFileDownload` power the first-party server-detail file manager. The page renders a generic server-root entry, requests live listings through `files.list`, reads snapshots through `files.read`, saves through `files.write`, and stages browser uploads as server-instance artifacts before Run pulls input chunks on the dedicated file-transfer channel. Server detail may call only these server-file APIs plus the encapsulated download helper; it must not call raw artifact-transfer methods directly. +- `getServerFileWorkspace`, `browseServerFiles`, `listServerFiles`, `refreshServerFiles`, `readServerFile`, `getServerFileReadSnapshot`, `writeServerFile`, `uploadServerFile`, and `prepareServerFileDownload` power the first-party server-detail file manager. The page renders a generic server-root entry and uses `browseServerFiles` as the live directory path for open, refresh, and search actions; compatibility `list`/`refresh` clients remain available for older flows. File reads use `files.read`, saves use `files.write`, and browser uploads stage server-instance artifacts before Run pulls input chunks on the dedicated file-transfer channel. Server detail may call only these server-file APIs plus the encapsulated download helper; it must not call raw artifact-transfer methods directly. - `listArtifacts`, `openArtifactDownload`, and `readArtifactContent` use platform artifact routes for available job/server artifacts. Browser reads are chunked through `/artifacts/{id}/content` and must render only safe filenames, checksums, progress, and platform storage behavior. - `authorizePluginBridge` posts `PluginBridgeAuthorizeRequest` to `/plugin-bridge/authorize` for preflight decisions. - `executePluginBridge` posts `PluginBridgeExecuteRequest` to `/plugin-bridge/execute` from host-owned bridge dispatch utilities only. Plugin pages receive typed `PluginBridgeExecuteResponse` envelopes and never receive the platform API client, bearer token, raw provider key, run socket, host path, or storage credential. diff --git a/platform_web/pages/ServerDetailPage.test.tsx b/platform_web/pages/ServerDetailPage.test.tsx index 3546d62..64af2be 100644 --- a/platform_web/pages/ServerDetailPage.test.tsx +++ b/platform_web/pages/ServerDetailPage.test.tsx @@ -97,10 +97,13 @@ describe("ServerDetailPage config write approval", () => { it("keeps the server file manager list-first without plugin declaration gates", () => { expect(serverDetailPageSource).toContain("server-file-editor-overlay"); - expect(serverDetailPageSource).toContain("refreshRuntimeList({ silent: true })"); - expect(serverDetailPageSource).toContain("window.setInterval(() => void loadList({ silent: true }), 300)"); + expect(serverDetailPageSource).toContain("browseServerFiles"); + expect(serverDetailPageSource).toContain("browseList({ forceNew: true })"); + expect(serverDetailPageSource).toContain("window.setTimeout(() => void browseList({ silent: true }), 500)"); expect(serverDetailPageSource).toContain("if (!options.silent) setList({ status: \"loading\" })"); expect(serverDetailPageSource).toContain("serverFileListPendingLabel"); + expect(serverDetailPageSource).not.toContain("refreshRuntimeList"); + expect(serverDetailPageSource).not.toContain("window.setInterval(() => void loadList"); expect(serverDetailPageSource).not.toContain("server-file-layout"); expect(serverDetailPageSource).not.toContain("未声明目录"); expect(serverDetailPageSource).not.toContain("插件尚未声明"); diff --git a/platform_web/pages/ServerDetailPage.tsx b/platform_web/pages/ServerDetailPage.tsx index 0806f92..3019c53 100644 --- a/platform_web/pages/ServerDetailPage.tsx +++ b/platform_web/pages/ServerDetailPage.tsx @@ -563,14 +563,14 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio const [panelResult, setPanelResult] = useState<{ status: "pending" | "succeeded" | "failed"; label: string } | null>(null); const [uploadBusy, setUploadBusy] = useState(false); const [editor, setEditor] = useState({ entry: null, key: "", draft: "", loading: false, saving: false }); - const initialRuntimeLoadRef = useRef(false); + const browseRequestRef = useRef<{ key: string; idempotencyKey: string } | null>(null); const activeDirectory = workspace.status === "ready" ? workspace.data.directories.find((item) => item.key === directoryKey) : undefined; const canUpload = workspace.status === "ready" && Boolean(activeDirectory) && !uploadBusy; const entries = list.status === "ready" ? list.data.entries : []; const loadWorkspace = useCallback(async () => { - initialRuntimeLoadRef.current = false; + browseRequestRef.current = null; setWorkspace({ status: "loading" }); try { const response = await platformApiClient.getServerFileWorkspace(instance.id); @@ -584,56 +584,45 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio } }, [instance.id]); - const loadList = useCallback(async (options: { silent?: boolean } = {}) => { - if (!directoryKey) return; + const browseList = useCallback(async (options: { silent?: boolean; forceNew?: boolean; manual?: boolean } = {}): Promise => { + if (!directoryKey) return undefined; if (!options.silent) setList({ status: "loading" }); + const browseKey = `${directoryKey}:${relativePath}:${searchQuery}:${recursive}`; + if (options.forceNew || browseRequestRef.current?.key !== browseKey) { + browseRequestRef.current = { key: browseKey, idempotencyKey: serverFileIdempotency("browse", instance.id, `${browseKey}:${Date.now()}`) }; + } + const operationId = options.manual ? operations.begin({ intent: "读取文件目录", targetKind: "server", targetId: instance.id, requester: session.displayName }) : ""; + if (!options.silent) setPanelResult({ status: "pending", label: "正在读取目录…" }); try { - const response = await platformApiClient.listServerFiles(instance.id, { directoryKey, path: relativePath || undefined, query: searchQuery || undefined, recursive }); + const response = await platformApiClient.browseServerFiles(instance.id, { directoryKey, path: relativePath || undefined, query: searchQuery || undefined, recursive, idempotencyKey: browseRequestRef.current?.idempotencyKey }); setList({ status: "ready", data: response }); - if (options.silent && response.state === "ready") setPanelResult(null); - if (options.silent && response.state === "pending") setPanelResult({ status: "pending", label: serverFileListPendingLabel(response) }); - if (options.silent && response.state === "failed") setPanelResult({ status: "failed", label: response.reason ?? "目录刷新失败" }); + if (operationId) operations.succeed(operationId, response.state === "ready" ? "目录读取完成" : "正在读取目录", response.job); + if (response.state === "ready") setPanelResult(null); + else if (response.state === "failed") setPanelResult({ status: "failed", label: response.reason ?? "目录读取失败" }); + else setPanelResult({ status: "pending", label: serverFileListPendingLabel(response) }); + return response; } catch (error) { setList({ status: "error", reason: error instanceof Error ? error.message : "文件列表加载失败" }); + if (operationId) operations.fail(operationId, error instanceof Error ? error.message : "目录读取失败", operationId); + setPanelResult({ status: "failed", label: error instanceof Error ? error.message : "目录读取失败" }); + return undefined; } - }, [directoryKey, instance.id, recursive, relativePath, searchQuery]); + }, [directoryKey, instance.id, operations, recursive, relativePath, searchQuery, session.displayName]); useEffect(() => { void loadWorkspace(); }, [loadWorkspace]); - const refreshRuntimeList = useCallback(async (options: { silent?: boolean } = {}) => { - if (!directoryKey) return; - const operationId = options.silent ? "" : operations.begin({ intent: "刷新文件目录", targetKind: "server", targetId: instance.id, requester: session.displayName }); - if (!options.silent) setPanelResult({ status: "pending", label: "正在向 Run 请求实时目录…" }); - try { - const response = await platformApiClient.refreshServerFiles(instance.id, { directoryKey, path: relativePath || undefined, query: searchQuery || undefined, recursive, idempotencyKey: serverFileIdempotency("list", instance.id, directoryKey) }); - setList({ status: "ready", data: response }); - if (operationId) operations.succeed(operationId, `目录刷新任务 ${response.job?.id ?? "已派发"}`, response.job); - if (response.state === "ready") setPanelResult(null); - else setPanelResult({ status: serverFileListResultStatus(response), label: serverFileListResultLabel(response) }); - } catch (error) { - const reason = error instanceof Error ? error.message : "目录刷新失败"; - if (operationId) operations.fail(operationId, reason, operationId); - setPanelResult({ status: "failed", label: reason }); - } - }, [directoryKey, instance.id, operations, relativePath, recursive, searchQuery, session.displayName]); - useEffect(() => { if (workspace.status !== "ready" || !directoryKey) return; - if (!initialRuntimeLoadRef.current) { - initialRuntimeLoadRef.current = true; - void refreshRuntimeList({ silent: true }); - return; - } - void loadList(); - }, [directoryKey, loadList, refreshRuntimeList, workspace.status]); + void browseList({ forceNew: true }); + }, [browseList, directoryKey, relativePath, recursive, searchQuery, workspace.status]); useEffect(() => { if (workspace.status !== "ready" || !directoryKey || list.status !== "ready" || list.data.state !== "pending") return; - const timer = window.setInterval(() => void loadList({ silent: true }), 300); - return () => window.clearInterval(timer); - }, [directoryKey, list, loadList, workspace.status]); + const timer = window.setTimeout(() => void browseList({ silent: true }), 500); + return () => window.clearTimeout(timer); + }, [browseList, directoryKey, list, workspace.status]); async function openEntry(entry: ServerFileEntryResponse) { if (entry.kind === "directory") { @@ -662,7 +651,7 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio const dispatch = await platformApiClient.readServerFile(instance.id, { key, idempotencyKey: serverFileIdempotency("read", instance.id, key) }); operations.succeed(operationId, `读取任务 ${dispatch.job.id} 已派发`, dispatch.job); setEditor({ entry, key, draft: "", snapshot, loading: false, saving: false, message: snapshot.reason ?? "读取任务已派发;Run 返回后再次打开即可编辑。" }); - setPanelResult({ status: "pending", label: `读取任务已派发:${dispatch.job.id}` }); + setPanelResult({ status: "pending", label: "正在读取文件…" }); } catch (error) { const reason = error instanceof Error ? error.message : "文件读取失败"; setEditor({ entry, key, draft: "", loading: false, saving: false, error: reason }); @@ -678,8 +667,8 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio const dispatch = await platformApiClient.writeServerFile(instance.id, { key: editor.key, content: editor.draft, expectedVersion: editor.snapshot?.version, expectedChecksum: editor.snapshot?.checksum, idempotencyKey: serverFileIdempotency("write", instance.id, editor.key) }); operations.succeed(operationId, `写入任务 ${dispatch.job.id} 已派发`, dispatch.job); setEditor((current) => ({ ...current, saving: false, message: "保存任务已派发;Run 会在工作区内原子写入。" })); - setPanelResult({ status: "pending", label: `写入任务已派发:${dispatch.job.id}` }); - await loadList(); + setPanelResult({ status: "pending", label: "正在保存文件…" }); + await browseList({ forceNew: true }); } catch (error) { const reason = error instanceof Error ? error.message : "文件保存失败"; operations.fail(operationId, reason, operationId); @@ -692,12 +681,19 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio const key = serverFileEntryKey(entry); if (!key) return; const operationId = operations.begin({ intent: "下载文件", targetKind: "server", targetId: instance.id, requester: session.displayName }); + const idempotencyKey = serverFileIdempotency("download", instance.id, key); setPanelResult({ status: "pending", label: "正在准备文件下载…" }); try { - const result = await platformApiClient.prepareServerFileDownload(instance.id, { key, idempotencyKey: serverFileIdempotency("download", instance.id, key) }); + let result = await platformApiClient.prepareServerFileDownload(instance.id, { key, idempotencyKey }); + for (let attempt = 0; result.status === "pending" && attempt < 240; attempt += 1) { + setPanelResult({ status: "pending", label: `正在准备下载… ${Math.min(99, Math.max(1, attempt))}%` }); + await new Promise((resolve) => window.setTimeout(resolve, 500)); + result = await platformApiClient.prepareServerFileDownload(instance.id, { key, idempotencyKey }); + } + if (result.status !== "ready") throw new Error(result.reason ?? "文件下载准备超时"); const message = await downloadServerFileResult(platformApiClient, result); operations.succeed(operationId, message, result.job); - setPanelResult({ status: result.status === "ready" ? "succeeded" : "pending", label: message }); + setPanelResult({ status: "succeeded", label: message }); } catch (error) { const reason = error instanceof Error ? error.message : "文件下载失败"; operations.fail(operationId, reason, operationId); @@ -719,8 +715,8 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio try { const response = await platformApiClient.uploadServerFile(instance.id, { directoryKey, relativePath: relativePath || undefined, file, idempotencyKey: serverFileIdempotency("upload", instance.id, file.name) }); operations.succeed(operationId, `上传已暂存,写入任务 ${response.job.id} 已派发`, response.job); - setPanelResult({ status: "pending", label: `上传已走独立文件通道排队:${response.relativePath}` }); - await loadList(); + setPanelResult({ status: "pending", label: `上传已提交:${response.relativePath}` }); + await browseList({ forceNew: true }); } catch (error) { const reason = error instanceof Error ? error.message : "文件上传失败"; operations.fail(operationId, reason, operationId); @@ -777,7 +773,7 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio
- + @@ -785,7 +781,7 @@ function ServerFilesSection({ instance, session, operations }: ServerFilesSectio
{panelResult && } {list.status === "loading" && } - {list.status === "error" && void loadList()} compact />} + {list.status === "error" && void browseList({ forceNew: true })} compact />} {list.status === "ready" && (
@@ -848,30 +844,14 @@ function serverFileEntryRowKey(entry: ServerFileEntryResponse): string { return `${entry.kind}:${entry.directoryKey}:${entry.relativePath ?? ""}:${entry.logicalKey ?? ""}:${entry.name}`; } -function serverFileListResultStatus(response: ServerFileListResponse): "pending" | "succeeded" | "failed" { - if (response.state === "ready") return "succeeded"; - if (response.state === "failed") return "failed"; - return "pending"; -} - -function serverFileListResultLabel(response: ServerFileListResponse): string { - if (response.state === "failed") return response.reason ?? "目录刷新失败"; - return serverFileListPendingLabel(response); -} - function serverFileListEmptyLabel(response: ServerFileListResponse): string { if (response.state === "pending") return serverFileListPendingLabel(response); - if (response.state === "failed") return response.reason ?? "目录刷新失败"; + if (response.state === "failed") return response.reason ?? "目录读取失败"; return "当前目录暂无文件。"; } function serverFileListPendingLabel(response: ServerFileListResponse): string { - const job = response.job; - if (!job) return response.reason ?? "正在读取当前目录,Run 返回后会自动更新。"; - const progress = job.progress?.message?.trim(); - const attempt = job.attempt > 0 ? ` · 第 ${job.attempt} 次尝试` : ""; - const nextAttempt = job.state === "retrying" && job.nextAttemptAt ? ` · 下次 ${formatDateTime(job.nextAttemptAt)}` : ""; - return `文件刷新任务 ${serverFileJobStateLabel(job.state)}${attempt}${nextAttempt}${progress ? ` · ${progress}` : ""}`; + return "正在读取当前目录…"; } function serverFileJobStateLabel(state: JobResponse["state"]): string { diff --git a/platform_web/utils/serverFileTransfer.ts b/platform_web/utils/serverFileTransfer.ts index 80c0bb2..9aac411 100644 --- a/platform_web/utils/serverFileTransfer.ts +++ b/platform_web/utils/serverFileTransfer.ts @@ -12,19 +12,49 @@ export async function downloadServerFileResult(client: PlatformApiClient, result if (!result.artifact) { return "文件内容尚未可用。"; } - const chunks: ArrayBuffer[] = []; + const writer = await createFileWriter(result.filename || result.artifact.filename); let offset = 0; const chunkSize = Math.max(1, result.artifact.chunkSizeBytes || 1024 * 1024); - while (offset < result.artifact.sizeBytes) { - const chunk = await client.readArtifactContent(result.artifact.artifactId, offset, Math.min(chunkSize, result.artifact.sizeBytes - offset)); - chunks.push(chunk.payload); - offset += chunk.contentLength; - if (chunk.contentLength <= 0) break; + try { + while (offset < result.artifact.sizeBytes) { + const chunk = await client.readArtifactContent(result.artifact.artifactId, offset, Math.min(chunkSize, result.artifact.sizeBytes - offset)); + if (chunk.contentLength <= 0) throw new Error("文件下载返回空数据"); + await writer.write(chunk.payload); + offset += chunk.contentLength; + } + await writer.close(); + } catch (error) { + await writer.abort(); + throw error; } - saveBlob(new Blob(chunks, { type: result.artifact.contentType || "application/octet-stream" }), result.filename || result.artifact.filename); return `下载已开始:${result.filename || result.artifact.filename}`; } +interface FileWriter { + write(data: ArrayBuffer): Promise; + close(): Promise; + abort(): Promise; +} + +async function createFileWriter(filename: string): Promise { + const picker = (window as Window & { showSaveFilePicker?: (options?: { suggestedName?: string }) => Promise<{ createWritable(): Promise }> }).showSaveFilePicker; + if (picker) { + const handle = await picker({ suggestedName: safeFilename(filename) }); + const writable = await handle.createWritable(); + return { + write: async (data) => { await writable.write(data); }, + close: async () => { await writable.close(); }, + abort: async () => { await writable.abort(); } + }; + } + const chunks: ArrayBuffer[] = []; + return { + write: async (data) => { chunks.push(data); }, + close: async () => { saveBlob(new Blob(chunks, { type: "application/octet-stream" }), safeFilename(filename)); }, + abort: async () => { chunks.length = 0; } + }; +} + function saveBlob(blob: Blob, filename: string) { const url = URL.createObjectURL(blob); const anchor = document.createElement("a");