145 lines
7.2 KiB
Go
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)
|
|
}
|
|
}
|