Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/codealmanac/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -344,6 +344,7 @@ def create_workflows(
sync = SyncWorkflow(
services.repositories,
services.sources,
services.runs,
queue,
SyncStateStore(services.local_state.database_path),
)
Expand Down
8 changes: 8 additions & 0 deletions src/codealmanac/services/runs/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from datetime import UTC, datetime

from codealmanac.core.errors import NotFoundError
from codealmanac.services.harnesses.models import HarnessTranscriptRef
from codealmanac.services.repositories.models import Repository, RepositoryName
from codealmanac.services.repositories.requests import SelectRepositoryRequest
from codealmanac.services.repositories.service import RepositoriesService
Expand Down Expand Up @@ -75,6 +76,13 @@ def list(self, request: ListRunsRequest) -> tuple[RunRecord, ...]:
repository_id = None if repository is None else repository.repository_id
return self.store.list(request.limit, repository_id=repository_id)

def list_harness_transcripts(
self,
repository_id: str,
) -> tuple[HarnessTranscriptRef, ...]:
self.repositories.get(repository_id)
return self.store.list_harness_transcripts(repository_id)

def show(self, request: ShowRunRequest) -> RunRecord:
record = self.store.read(request.run_id)
self.require_run_matches_repository(record, request.repository_name)
Expand Down
11 changes: 11 additions & 0 deletions src/codealmanac/services/runs/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,17 @@ def list(
with self.connect() as connection:
return list_run_records(connection, limit, repository_id)

def list_harness_transcripts(
self,
repository_id: str,
) -> tuple[HarnessTranscriptRef, ...]:
records = self.list(limit=None, repository_id=repository_id)
return tuple(
record.harness_transcript
for record in records
if record.harness_transcript is not None
)

def read(self, run_id: str) -> RunRecord:
with self.connect() as connection:
record = read_run_record(connection, run_id)
Expand Down
21 changes: 21 additions & 0 deletions src/codealmanac/workflows/sync/evaluation.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from codealmanac.services.repositories.models import Repository
from codealmanac.services.repositories.requests import SelectRepositoryRequest
from codealmanac.services.repositories.service import RepositoriesService
from codealmanac.services.runs.service import RunsService
from codealmanac.services.sources.models import TranscriptCandidate
from codealmanac.services.sources.requests import DiscoverTranscriptsRequest
from codealmanac.services.sources.service import SourcesService
Expand All @@ -29,10 +30,12 @@ def __init__(
self,
repositories: RepositoriesService,
sources: SourcesService,
runs: RunsService,
state_store: SyncStateStore,
):
self.repositories = repositories
self.sources = sources
self.runs = runs
self.state_store = state_store

def evaluate(
Expand Down Expand Up @@ -62,6 +65,10 @@ def evaluate(
normalize_path(repository.root_path): repository
for repository in selected_repositories
}
owned_sessions_by_repository_id = {
repository.repository_id: self.owned_sessions(repository)
for repository in selected_repositories
}
for transcript in transcripts:
repository = repositories_by_path.get(normalize_path(transcript.cwd))
if repository is None:
Expand All @@ -70,6 +77,12 @@ def evaluate(
if transcript.modified_at < active_since:
skipped.append(skipped_transcript(transcript, "inactive"))
continue
if (
transcript.app.value,
transcript.session_id,
) in owned_sessions_by_repository_id[repository.repository_id]:
skipped.append(skipped_transcript(transcript, "codealmanac-run"))
continue
transcripts_by_repository_id[repository.repository_id].append(transcript)

repository_ingests = tuple(
Expand Down Expand Up @@ -109,6 +122,14 @@ def selected_repositories(
)
return (repository,)

def owned_sessions(self, repository: Repository) -> frozenset[tuple[str, str]]:
return frozenset(
(transcript.kind.value, transcript.session_id)
for transcript in self.runs.list_harness_transcripts(
repository.repository_id
)
)


def transcript_sort_key(candidate: TranscriptCandidate) -> tuple[str, str, str]:
return (
Expand Down
3 changes: 3 additions & 0 deletions src/codealmanac/workflows/sync/service.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
from datetime import UTC, datetime

from codealmanac.services.repositories.service import RepositoriesService
from codealmanac.services.runs.service import RunsService
from codealmanac.services.sources.service import SourcesService
from codealmanac.workflows.run_queue.service import RunQueue
from codealmanac.workflows.sync.evaluation import SyncEvaluator
Expand All @@ -23,12 +24,14 @@ def __init__(
self,
repositories: RepositoriesService,
sources: SourcesService,
runs: RunsService,
queue: RunQueue,
state_store: SyncStateStore,
):
self.evaluator = SyncEvaluator(
repositories=repositories,
sources=sources,
runs=runs,
state_store=state_store,
)
self.executor = SyncIngestQueue(
Expand Down
92 changes: 85 additions & 7 deletions tests/test_sync_workflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,18 @@
from datetime import UTC, datetime, timedelta
from pathlib import Path

import pytest
from conftest import initialize_repository

from codealmanac.app import create_app
from codealmanac.services.harnesses.models import HarnessKind
from codealmanac.services.harnesses.models import HarnessKind, HarnessTranscriptRef
from codealmanac.services.runs.models import RunKind, RunWorkerSpawnResult
from codealmanac.services.runs.requests import ReadRunSpecRequest, SpawnRunWorkerRequest
from codealmanac.services.runs.requests import (
ReadRunSpecRequest,
RecordRunHarnessTranscriptRequest,
SpawnRunWorkerRequest,
StartRunRequest,
)
from codealmanac.services.sources.models import TranscriptApp, TranscriptCandidate
from codealmanac.services.sources.requests import DiscoverTranscriptsRequest
from codealmanac.settings import AppConfig
Expand All @@ -16,9 +22,12 @@


class FakeTranscriptDiscoveryAdapter:
app = TranscriptApp.CODEX

def __init__(self, candidates: tuple[TranscriptCandidate, ...]):
def __init__(
self,
candidates: tuple[TranscriptCandidate, ...],
app: TranscriptApp = TranscriptApp.CODEX,
):
self.app = app
self.candidates = candidates
self.requests: list[DiscoverTranscriptsRequest] = []

Expand Down Expand Up @@ -104,6 +113,73 @@ def test_sync_uses_exact_registered_cwd_without_root_hopping(
assert summary.skipped[0].reason == "unregistered-cwd"


@pytest.mark.parametrize(
("transcript_app", "harness_kind"),
(
(TranscriptApp.CLAUDE, HarnessKind.CLAUDE),
(TranscriptApp.CODEX, HarnessKind.CODEX),
),
)
def test_sync_skips_transcripts_from_codealmanac_runs(
tmp_path: Path,
isolated_home: Path,
transcript_app: TranscriptApp,
harness_kind: HarnessKind,
):
repo = initialized_repo(tmp_path, "repo")
own_work = transcript_candidate(
repo,
tmp_path / "own-work.jsonl",
current_time(),
app=transcript_app,
session_id=f"{transcript_app.value}-codealmanac-session",
)
user_work = transcript_candidate(
repo,
tmp_path / "user-work.jsonl",
current_time(),
app=transcript_app,
session_id=f"{transcript_app.value}-user-session",
)
app, _, spawner = app_with_sync(
isolated_home,
candidates=(own_work, user_work),
app=transcript_app,
)
repository = initialize_repository(app, path=repo)
run = app.runs.start(
StartRunRequest(
repository_id=repository.repository_id,
kind=RunKind.INGEST,
title="Sync own work",
)
)
app.runs.record_harness_transcript(
RecordRunHarnessTranscriptRequest(
run_id=run.run_id,
transcript=HarnessTranscriptRef(
kind=harness_kind,
session_id=f"{transcript_app.value}-codealmanac-session",
),
)
)

summary = app.workflows.sync.status(
SyncStatusRequest(
apps=(transcript_app,),
now=current_time(),
)
)

assert summary.eligible == 1
assert summary.ready[0].transcript_paths == (user_work.transcript_path,)
assert summary.skipped[0].reason == "codealmanac-run"
assert summary.skipped[0].session_id == (
f"{transcript_app.value}-codealmanac-session"
)
assert spawner.requests == []


def test_sync_run_queues_one_ingest_run_per_repository_and_records_scan_time(
tmp_path: Path,
isolated_home: Path,
Expand Down Expand Up @@ -240,10 +316,11 @@ def app_with_sync(
isolated_home: Path,
*,
candidates: tuple[TranscriptCandidate, ...],
app: TranscriptApp = TranscriptApp.CODEX,
database_path: Path | None = None,
spawner: SyncWorkerSpawner | None = None,
):
adapter = FakeTranscriptDiscoveryAdapter(candidates)
adapter = FakeTranscriptDiscoveryAdapter(candidates, app=app)
selected_spawner = spawner or SyncWorkerSpawner()
app = create_app(
AppConfig(
Expand All @@ -261,11 +338,12 @@ def transcript_candidate(
path: Path,
modified_at: datetime,
*,
app: TranscriptApp = TranscriptApp.CODEX,
session_id: str = "session-1",
) -> TranscriptCandidate:
path.write_text("transcript\n", encoding="utf-8")
return TranscriptCandidate(
app=TranscriptApp.CODEX,
app=app,
session_id=session_id,
transcript_path=path,
cwd=repo,
Expand Down