Skip to content

Commit 6fa7ffa

Browse files
authored
feat(artifact): import normalized peer checkpoints (#1299)
Artifact export currently establishes durable publication authority, but artifacts already present in a store cannot be consumed into the normalized archive with equivalent crash safety. That leaves round-trip synchronization one-way and makes later folder or service transports unsafe to add because inbound work has no durable claim, version-deferral, or landing boundary. This change adds exact checkpoint claims, monotonic peer heads, complete landing maps, and per-session provenance that survive full resync. Imports authenticate bounded artifact reads, accept semantically valid peer JSON under its stored identity, defer missing or future dependencies independently, quarantine deterministic corruption, and rewrite complete sessions under origin-qualified identities without carrying source-machine or unverified secret state. The coordinator acknowledges a checkpoint only after its available session closures and exact landing map are durable. Retries converge across each write boundary, unchanged large checkpoints perform closure work only for changed sessions, and repeated observations do not duplicate normalized messages or usage. Checkpoint absence remains non-destructive. Metadata replay, artifact transport, raw provider-source stewardship, and provider-file eviction remain separate follow-on scopes. Co-authored-by: Wes McKinney <wesm@users.noreply.github.com>
1 parent e7aa95f commit 6fa7ffa

21 files changed

Lines changed: 9584 additions & 53 deletions

internal/artifact/export.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import (
1313
"io/fs"
1414
"os"
1515
"sort"
16+
"strings"
1617

1718
"go.kenn.io/agentsview/internal/db"
1819
)
@@ -407,6 +408,9 @@ func exportFullToStoreWithDrainRoundsAndLimits(
407408
if _, _, err := exportClaimedSessionToStore(
408409
ctx, database, store, origin, sess, limits,
409410
); err != nil {
411+
if isDeterministicArtifactExportError(err) {
412+
continue
413+
}
410414
return result, err
411415
}
412416
}
@@ -502,6 +506,12 @@ func exportLoadedSessionToStore(
502506
usageEvents []db.UsageEvent,
503507
limits artifactLimits,
504508
) (string, bool, error) {
509+
if sess.ID == "" || strings.Contains(sess.ID, "~") {
510+
return "", false, rejectArtifactExportf(
511+
"native session ID %q is invalid: must be non-empty and exclude '~'",
512+
sess.ID,
513+
)
514+
}
505515
if len(messages) > limits.sessionMessages {
506516
return "", false, rejectArtifactExportf(
507517
"session message limit exceeded for %s: got %d, limit %d",

internal/artifact/export_limits_test.go

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -112,6 +112,46 @@ func TestExportClassifiesBoundedLoadLimit(t *testing.T) {
112112
assertNoPublishedArtifacts(t, store, contractOrigin)
113113
}
114114

115+
func TestExportRejectsNativeSessionIDsImporterCannotRepresent(t *testing.T) {
116+
tests := []struct {
117+
name string
118+
sessionID string
119+
}{
120+
{name: "empty", sessionID: ""},
121+
{name: "reserved separator", sessionID: "native~session"},
122+
}
123+
for _, tt := range tests {
124+
t.Run(tt.name, func(t *testing.T) {
125+
store := newTestArtifactStore(t)
126+
session := &db.Session{
127+
ID: tt.sessionID,
128+
Machine: "local",
129+
MessageCount: 1,
130+
UserMessageCount: 1,
131+
}
132+
messages := []db.Message{{
133+
SessionID: tt.sessionID,
134+
Ordinal: 0,
135+
Role: "user",
136+
Content: "hello",
137+
}}
138+
139+
_, _, err := exportLoadedSessionToStore(
140+
t.Context(),
141+
store,
142+
contractOrigin,
143+
session,
144+
messages,
145+
nil,
146+
productionArtifactLimits(),
147+
)
148+
require.ErrorIs(t, err, ErrArtifactExportRejected)
149+
assert.Contains(t, err.Error(), "native session ID")
150+
assertNoPublishedArtifacts(t, store, contractOrigin)
151+
})
152+
}
153+
}
154+
115155
func TestExportRejectsNestedAmplificationBeforePublication(t *testing.T) {
116156
tests := []struct {
117157
name string

internal/artifact/export_test.go

Lines changed: 102 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -468,6 +468,108 @@ func TestExportContinuesPastInvalidManifestValue(t *testing.T) {
468468
assert.Contains(t, rejection.Error, "encoding manifest for invalid")
469469
}
470470

471+
func TestFullExportSkipsPreviouslyRejectedNativeSessionID(t *testing.T) {
472+
database := testExportDB(t)
473+
store := newTestArtifactStore(t)
474+
seedSession(t, database, "invalid~native", "alpha")
475+
476+
rejected, err := ExportToStore(
477+
t.Context(), database, store, ExportOptions{Origin: contractOrigin},
478+
)
479+
require.NoError(t, err)
480+
require.Equal(t, 1, rejected.RejectedSessions)
481+
482+
result, err := ExportToStore(
483+
t.Context(), database, store, ExportOptions{
484+
Origin: contractOrigin,
485+
Full: true,
486+
},
487+
)
488+
require.NoError(t, err)
489+
assert.Zero(t, result.ExportedSessions)
490+
assert.Zero(t, result.RejectedSessions)
491+
pending, err := database.PendingArtifactExports(t.Context(), 10)
492+
require.NoError(t, err)
493+
assert.Empty(t, pending)
494+
}
495+
496+
func TestFullExportRemovesPublishedNativeSessionIDRejectedByCurrentWire(t *testing.T) {
497+
ctx := t.Context()
498+
databasePath := filepath.Join(t.TempDir(), "archive.db")
499+
database, err := db.Open(databasePath)
500+
require.NoError(t, err)
501+
require.NoError(t, database.SetSyncState(originStateKey, contractOrigin))
502+
store := newTestArtifactStore(t)
503+
seedSession(t, database, "legacy~native", "alpha")
504+
505+
claims, err := database.ArtifactExportClaims(ctx, []string{"legacy~native"})
506+
require.NoError(t, err)
507+
require.Len(t, claims, 1)
508+
manifestHash := strings.Repeat("a", 64)
509+
revision, changed, err := database.ApplyArtifactPublicationChanges(
510+
ctx, contractOrigin, []db.ArtifactPublicationChange{{
511+
SessionID: "legacy~native",
512+
Generation: claims[0].Generation,
513+
ManifestHash: manifestHash,
514+
SourceFingerprint: manifestHash,
515+
}},
516+
)
517+
require.NoError(t, err)
518+
require.True(t, changed)
519+
publications, mapDigest, mapRevision, err := spoolArtifactPublicationMap(
520+
ctx, database, contractOrigin,
521+
)
522+
require.NoError(t, err)
523+
require.Equal(t, revision, mapRevision)
524+
checkpointBody, checkpointIdentity, err := spoolArtifactCheckpoint(
525+
ctx, publications, contractOrigin, 1,
526+
)
527+
require.NoError(t, err)
528+
checkpointRef := requireContractRef(
529+
t, contractOrigin, KindCheckpoints, "cp-0000000001.json",
530+
)
531+
_, err = store.Create(
532+
ctx, checkpointRef, checkpointIdentity,
533+
canonicalArtifactMediaType(KindCheckpoints), checkpointBody,
534+
)
535+
require.NoError(t, err)
536+
require.NoError(t, closeAndRemoveExportSpool(publications))
537+
require.NoError(t, closeAndRemoveExportSpool(checkpointBody))
538+
require.NoError(t, database.RecordArtifactCheckpointHead(
539+
ctx, db.ArtifactCheckpointHead{
540+
Origin: contractOrigin,
541+
Sequence: 1,
542+
PublicationRevision: revision,
543+
SessionMapSHA256: mapDigest,
544+
CheckpointSHA256: checkpointIdentity.SHA256,
545+
CheckpointSize: checkpointIdentity.Size,
546+
},
547+
claims,
548+
))
549+
require.NoError(t, database.Close())
550+
551+
database, err = db.Open(databasePath)
552+
require.NoError(t, err)
553+
t.Cleanup(func() { require.NoError(t, database.Close()) })
554+
result, err := ExportToStore(ctx, database, store, ExportOptions{
555+
Origin: contractOrigin,
556+
Full: true,
557+
})
558+
require.NoError(t, err)
559+
assert.Equal(t, 1, result.RejectedSessions)
560+
published := latestStoreCheckpointForTest(t, store, contractOrigin)
561+
assert.NotContains(t, published.Sessions, contractOrigin+"~legacy~native")
562+
pending, err := database.PendingArtifactExports(ctx, 10)
563+
require.NoError(t, err)
564+
assert.Empty(t, pending)
565+
rejection, ok, err := database.GetArtifactExportRejection(
566+
ctx, "legacy~native",
567+
)
568+
require.NoError(t, err)
569+
require.True(t, ok)
570+
assert.Contains(t, rejection.Error, "native session ID")
571+
}
572+
471573
func TestExportTransientFailureDoesNotReject(t *testing.T) {
472574
database := testExportDB(t)
473575
seedSession(t, database, "rejected", "alpha")

0 commit comments

Comments
 (0)