Implement Apex Agent MVP scaffold
This commit is contained in:
@@ -0,0 +1,2 @@
|
||||
"""Local artifact and task metadata storage."""
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from dataclasses import asdict, is_dataclass
|
||||
from datetime import date, datetime
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
|
||||
class ArtifactStore:
|
||||
def __init__(self, root: str | Path = "tasks/history") -> None:
|
||||
self.root = Path(root)
|
||||
|
||||
def task_dir(self, task_id: str) -> Path:
|
||||
return self.root / task_id
|
||||
|
||||
def write_step(
|
||||
self,
|
||||
*,
|
||||
task_id: str,
|
||||
index: int,
|
||||
screenshot: bytes | None,
|
||||
record: dict[str, Any],
|
||||
) -> dict[str, str | None]:
|
||||
task_dir = self.task_dir(task_id)
|
||||
task_dir.mkdir(parents=True, exist_ok=True)
|
||||
stem = f"{index:03d}"
|
||||
screenshot_path: Path | None = None
|
||||
if screenshot is not None:
|
||||
screenshot_path = task_dir / f"{stem}.png"
|
||||
screenshot_path.write_bytes(screenshot)
|
||||
|
||||
json_path = task_dir / f"{stem}.json"
|
||||
payload = {
|
||||
**record,
|
||||
"screenshot_path": str(screenshot_path) if screenshot_path else None,
|
||||
}
|
||||
json_path.write_text(
|
||||
json.dumps(_jsonable(payload), ensure_ascii=False, indent=2),
|
||||
encoding="utf-8",
|
||||
)
|
||||
return {
|
||||
"json_path": str(json_path),
|
||||
"screenshot_path": str(screenshot_path) if screenshot_path else None,
|
||||
}
|
||||
|
||||
def read_steps(self, task_id: str) -> list[dict[str, Any]]:
|
||||
task_dir = self.task_dir(task_id)
|
||||
if not task_dir.exists():
|
||||
return []
|
||||
steps = []
|
||||
for path in sorted(task_dir.glob("*.json")):
|
||||
steps.append(json.loads(path.read_text(encoding="utf-8")))
|
||||
return steps
|
||||
|
||||
|
||||
def _jsonable(value: Any) -> Any:
|
||||
if hasattr(value, "to_dict"):
|
||||
return value.to_dict()
|
||||
if is_dataclass(value):
|
||||
return asdict(value)
|
||||
if isinstance(value, dict):
|
||||
return {key: _jsonable(inner) for key, inner in value.items()}
|
||||
if isinstance(value, list):
|
||||
return [_jsonable(inner) for inner in value]
|
||||
if isinstance(value, tuple):
|
||||
return [_jsonable(inner) for inner in value]
|
||||
if isinstance(value, (datetime, date)):
|
||||
return value.isoformat()
|
||||
return value
|
||||
|
||||
@@ -0,0 +1,97 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import sqlite3
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from core.models import Task, TaskStatus, utc_now
|
||||
|
||||
|
||||
class TaskMetadataStore:
|
||||
def __init__(self, db_path: str | Path = "tasks/tasks.sqlite3") -> None:
|
||||
self.db_path = Path(db_path)
|
||||
self.db_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
self._ensure_schema()
|
||||
|
||||
def create_task(self, task: Task) -> None:
|
||||
with self._connect() as connection:
|
||||
connection.execute(
|
||||
"""
|
||||
insert into tasks (
|
||||
id, goal, device_id, status, created_at, updated_at,
|
||||
completed_at, failure_reason
|
||||
) values (?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
task.id,
|
||||
task.goal,
|
||||
task.device_id,
|
||||
task.status,
|
||||
task.created_at.isoformat(),
|
||||
task.updated_at.isoformat(),
|
||||
task.completed_at.isoformat() if task.completed_at else None,
|
||||
task.failure_reason,
|
||||
),
|
||||
)
|
||||
|
||||
def update_task(
|
||||
self,
|
||||
task_id: str,
|
||||
*,
|
||||
status: TaskStatus | None = None,
|
||||
failure_reason: str | None = None,
|
||||
completed: bool = False,
|
||||
) -> None:
|
||||
updates: dict[str, Any] = {"updated_at": utc_now().isoformat()}
|
||||
if status:
|
||||
updates["status"] = status
|
||||
if failure_reason is not None:
|
||||
updates["failure_reason"] = failure_reason
|
||||
if completed:
|
||||
updates["completed_at"] = utc_now().isoformat()
|
||||
|
||||
assignments = ", ".join(f"{key} = ?" for key in updates)
|
||||
values = [*updates.values(), task_id]
|
||||
with self._connect() as connection:
|
||||
connection.execute(
|
||||
f"update tasks set {assignments} where id = ?",
|
||||
values,
|
||||
)
|
||||
|
||||
def get_task(self, task_id: str) -> dict[str, Any] | None:
|
||||
with self._connect() as connection:
|
||||
row = connection.execute(
|
||||
"select * from tasks where id = ?",
|
||||
(task_id,),
|
||||
).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
def list_tasks(self) -> list[dict[str, Any]]:
|
||||
with self._connect() as connection:
|
||||
rows = connection.execute(
|
||||
"select * from tasks order by created_at desc"
|
||||
).fetchall()
|
||||
return [dict(row) for row in rows]
|
||||
|
||||
def _ensure_schema(self) -> None:
|
||||
with self._connect() as connection:
|
||||
connection.execute(
|
||||
"""
|
||||
create table if not exists tasks (
|
||||
id text primary key,
|
||||
goal text not null,
|
||||
device_id text not null,
|
||||
status text not null,
|
||||
created_at text not null,
|
||||
updated_at text not null,
|
||||
completed_at text,
|
||||
failure_reason text
|
||||
)
|
||||
"""
|
||||
)
|
||||
|
||||
def _connect(self) -> sqlite3.Connection:
|
||||
connection = sqlite3.connect(self.db_path)
|
||||
connection.row_factory = sqlite3.Row
|
||||
return connection
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from storage.artifact_store import ArtifactStore
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TimelineRecord:
|
||||
index: int
|
||||
scene: dict[str, Any]
|
||||
prompt: str
|
||||
tool_call: dict[str, Any]
|
||||
result: dict[str, Any]
|
||||
timestamp: str
|
||||
screenshot_path: str | None = None
|
||||
|
||||
|
||||
class Timeline:
|
||||
def __init__(self, artifact_store: ArtifactStore | None = None) -> None:
|
||||
self.artifact_store = artifact_store or ArtifactStore()
|
||||
|
||||
def append(
|
||||
self,
|
||||
*,
|
||||
task_id: str,
|
||||
scene: Any,
|
||||
prompt: str,
|
||||
tool_call: dict[str, Any],
|
||||
result: dict[str, Any],
|
||||
screenshot: bytes | None = None,
|
||||
) -> TimelineRecord:
|
||||
index = len(self.read(task_id)) + 1
|
||||
record = {
|
||||
"index": index,
|
||||
"scene": scene,
|
||||
"prompt": prompt,
|
||||
"tool_call": tool_call,
|
||||
"result": result,
|
||||
"timestamp": datetime.now().astimezone().isoformat(),
|
||||
}
|
||||
paths = self.artifact_store.write_step(
|
||||
task_id=task_id,
|
||||
index=index,
|
||||
screenshot=screenshot,
|
||||
record=record,
|
||||
)
|
||||
return TimelineRecord(
|
||||
index=index,
|
||||
scene=record["scene"],
|
||||
prompt=prompt,
|
||||
tool_call=tool_call,
|
||||
result=result,
|
||||
timestamp=record["timestamp"],
|
||||
screenshot_path=paths["screenshot_path"],
|
||||
)
|
||||
|
||||
def read(self, task_id: str) -> list[dict[str, Any]]:
|
||||
return self.artifact_store.read_steps(task_id)
|
||||
|
||||
Reference in New Issue
Block a user