Files
browser/platform/service/durable_observability_test.go
2026-08-26 22:00:23 +08:00

145 lines
7.2 KiB
Go

package service
import (
"path/filepath"
"testing"
"time"
"browser.local/platform/domain"
"browser.local/platform/repo"
"browser.local/platform/validator"
)
func TestFileArtifactBodyStoreResumesTransferAfterServiceRestart(t *testing.T) {
root := t.TempDir()
metadata := filepath.Join(root, "metadata.json")
store, err := repo.NewFileStore(metadata)
if err != nil {
t.Fatalf("new file store: %v", err)
}
logStore, err := NewFileLogBodyStore(filepath.Join(root, "logs"))
if err != nil {
t.Fatalf("new log store: %v", err)
}
artifactStore, err := NewFileArtifactBodyStore(filepath.Join(root, "artifacts"))
if err != nil {
t.Fatalf("new artifact store: %v", err)
}
svc, err := NewCoreServiceWithDurableStores(store, logStore, artifactStore)
if err != nil {
t.Fatalf("new durable service: %v", err)
}
plugin, endpoint := createPluginAndRunEndpoint(t, svc)
if _, err := svc.CreateServerInstance(domain.ServerInstance{ID: "durable-server", PluginID: plugin.ID, RunEndpointID: endpoint.ID, Name: "Durable"}); err != nil {
t.Fatalf("create instance: %v", err)
}
if _, err := svc.CreateJob(domain.Job{ID: "durable-job", ServerInstanceID: "durable-server", RunEndpointID: endpoint.ID, Capability: "process.start", IdempotencyKey: "durable-job"}); err != nil {
t.Fatalf("create job: %v", err)
}
hello, err := svc.RegisterRunHello(validRunControlHello())
if err != nil {
t.Fatalf("register run: %v", err)
}
payload := []byte("durable transfer payload")
open, err := svc.OpenArtifactTransfer(domain.ArtifactTransferOpen{RunEndpointID: endpoint.ID, SessionToken: hello.SessionToken, ArtifactID: "durable-artifact", Direction: domain.ArtifactTransferDirectionUpload, OwnerKind: domain.ArtifactOwnerKindJob, OwnerID: "durable-job", SizeBytes: int64(len(payload)), ChunkSizeBytes: 8, Checksum: validator.BytesChecksum(payload), IdempotencyKey: "durable-transfer"})
if err != nil {
t.Fatalf("open transfer: %v", err)
}
first := validArtifactChunk(hello.SessionToken, open.TransferID, payload, 0, 8)
first.ArtifactID = "durable-artifact"
if _, err := svc.UploadArtifactChunk(first); err != nil {
t.Fatalf("upload first chunk: %v", err)
}
reloadedStore, err := repo.NewFileStore(metadata)
if err != nil {
t.Fatalf("reload metadata store: %v", err)
}
reloadedLogStore, err := NewFileLogBodyStore(filepath.Join(root, "logs"))
if err != nil {
t.Fatalf("reload log store: %v", err)
}
reloadedArtifacts, err := NewFileArtifactBodyStore(filepath.Join(root, "artifacts"))
if err != nil {
t.Fatalf("reload artifact store: %v", err)
}
restarted, err := NewCoreServiceWithDurableStores(reloadedStore, reloadedLogStore, reloadedArtifacts)
if err != nil {
t.Fatalf("restart service: %v", err)
}
status, err := restarted.QueryArtifactTransferStatus(domain.ArtifactTransferStatusQuery{RunEndpointID: endpoint.ID, SessionToken: hello.SessionToken, TransferID: open.TransferID, ArtifactID: "durable-artifact"})
if err != nil {
t.Fatalf("query resumed status: %v", err)
}
if status.NextMissingChunkIndex != 1 || len(status.ReceivedChunkIndexes) != 1 {
t.Fatalf("unexpected resumed transfer status: %+v", status)
}
for index := 1; index < open.TotalChunks; index++ {
chunk := validArtifactChunk(hello.SessionToken, open.TransferID, payload, index, 8)
chunk.ArtifactID = "durable-artifact"
if _, err := restarted.UploadArtifactChunk(chunk); err != nil {
t.Fatalf("upload resumed chunk %d: %v", index, err)
}
}
if _, err := restarted.CompleteArtifactTransfer(domain.ArtifactTransferComplete{RunEndpointID: endpoint.ID, SessionToken: hello.SessionToken, TransferID: open.TransferID, ArtifactID: "durable-artifact", Checksum: validator.BytesChecksum(payload), SizeBytes: int64(len(payload))}); err != nil {
t.Fatalf("complete resumed transfer: %v", err)
}
finalStore, _ := repo.NewFileStore(metadata)
finalArtifacts, _ := NewFileArtifactBodyStore(filepath.Join(root, "artifacts"))
finalService, err := NewCoreServiceWithDurableStores(finalStore, reloadedLogStore, finalArtifacts)
if err != nil {
t.Fatalf("final restart service: %v", err)
}
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)
}
}
func TestMetricsAndBackupsPersistWithRetentionRecovery(t *testing.T) {
svc := newTestCoreService()
plugin, endpoint := createPluginAndRunEndpoint(t, svc)
ownerSession := createServiceUserAndLogin(t, svc, domain.User{ID: "observability-owner", DisplayName: "Owner", Email: "observability-owner@example.test", Roles: []string{"server-owner"}, PasswordHash: "secret-password"})
instance, err := svc.CreateServerInstanceForSession(ownerSession, domain.ServerInstance{ID: "observability-server", PluginID: plugin.ID, RunEndpointID: endpoint.ID, Name: "Observability"})
if err != nil {
t.Fatalf("create instance: %v", err)
}
runHello := validRunControlHello()
runHello.RunEndpointID = endpoint.ID
registered, err := svc.RegisterRunHello(runHello)
if err != nil {
t.Fatalf("register run: %v", err)
}
collectedAt := time.Date(2026, 7, 18, 12, 0, 0, 0, time.UTC)
cpu := 42.0
if _, err := svc.IngestMetricBatch(domain.MetricBatchIngest{RunEndpointID: endpoint.ID, SessionToken: registered.SessionToken, Samples: []domain.MetricSample{{ServerInstanceID: instance.ID, CPUPercent: &cpu, Source: "run", CollectedAt: collectedAt}}}); err != nil {
t.Fatalf("ingest metrics: %v", err)
}
metrics, err := svc.ListMetricSamplesForSession(ownerSession, domain.MetricSampleFilter{ServerInstanceID: instance.ID, Limit: 10})
if err != nil || len(metrics) != 1 || metrics[0].CPUPercent == nil || *metrics[0].CPUPercent != cpu {
t.Fatalf("unexpected persisted metrics: %+v err=%v", metrics, err)
}
artifact, err := svc.CreateArtifact(domain.Artifact{ID: "backup-artifact", OwnerKind: domain.ArtifactOwnerKindServerInstance, OwnerID: instance.ID, SizeBytes: 12, Checksum: validator.BytesChecksum([]byte("backup bytes")), State: domain.ArtifactStateAvailable})
if err != nil {
t.Fatalf("create backup artifact: %v", err)
}
backup, err := svc.CreateBackupForSession(ownerSession, domain.BackupRecord{ID: "backup-1", ServerInstanceID: instance.ID, ArtifactID: artifact.ID})
if err != nil || backup.State != domain.BackupStatePending {
t.Fatalf("create backup record: %+v err=%v", backup, err)
}
if err := svc.RecoverIncompleteBackups(); err != nil {
t.Fatalf("recover backups: %v", err)
}
recovered, err := svc.GetBackupForSession(ownerSession, backup.ID)
if err != nil || recovered.State != domain.BackupStateFailed || recovered.RecoveryStatus == "" {
t.Fatalf("expected recoverable failed backup, record=%+v err=%v", recovered, err)
}
otherSession := createServiceUserAndLogin(t, svc, domain.User{ID: "observability-other", DisplayName: "Other", Email: "observability-other@example.test", Roles: []string{"server-owner"}, PasswordHash: "secret-password"})
if _, err := svc.ListMetricSamplesForSession(otherSession, domain.MetricSampleFilter{ServerInstanceID: instance.ID, Limit: 10}); err != ErrForbidden {
t.Fatalf("expected cross-owner metric denial, got %v", err)
}
if _, err := svc.GetBackupForSession(otherSession, backup.ID); err != ErrForbidden {
t.Fatalf("expected cross-owner backup denial, got %v", err)
}
}