Skip to content
Merged
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
10 changes: 6 additions & 4 deletions .github/workflows/plugin-api-v3.yml
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ jobs:
- uses: actions/checkout@v4
with:
repository: kachofugetsu09/akashic-agent
ref: 78e50d4dfb3f4348fff37d55d9c9bdd0e002164d
ref: 9da3a988a2bf62b0f550bd4f6bb98c4eeb1f56f5
path: .akashic-core
- uses: actions/setup-python@v5
with:
Expand All @@ -55,6 +55,7 @@ jobs:
- name: Verify Fitbit v3 composition
env:
AKASHIC_AGENT_ROOT: .akashic-core
AKASHIC_PLUGIN_FIXTURE_PYTHON: ${{ github.workspace }}/.venv/bin/python
PYTHONPATH: .akashic-core
run: .venv/bin/python -m pytest -q tests/
- uses: actions/setup-node@v4
Expand All @@ -70,11 +71,12 @@ jobs:
PYTHONPATH: .akashic-core
run: >-
.venv/bin/pyright --pythonpath .venv/bin/python --level error
plugin.py dashboard.py src/mcp_bridge.py
plugin.py dashboard.py src/content_adapter.py src/mcp_bridge.py
monitor/runtime_env.py
scripts/migrate_v2_data.py
tests/test_plugin.py tests/test_dashboard.py
tests/test_plugin.py tests/test_content_adapter.py tests/test_dashboard.py
tests/test_manager_integration.py tests/test_mcp_v3_runtime.py
tests/test_migrate_v2_data.py
tests/test_migrate_v2_data.py tests/test_runtime_env.py
- name: Compile Python sources
run: python -m compileall -q plugin.py dashboard.py src monitor scripts tests
- name: Check diff formatting
Expand Down
9 changes: 4 additions & 5 deletions akashic.plugin.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
schema_version = 1
name = "fitbit"
version = "3.0.0"
version = "3.1.0"
api_version = 3
entrypoint = "plugin.py"

Expand Down Expand Up @@ -37,10 +37,9 @@ startup_timeout_seconds = 15.0
name = "fitbit"
command = ["python", "run_mcp.py"]
required_tools = [
"get_proactive_events",
"get_sleep_context",
"acknowledge_events",
"fitbit_health_snapshot",
"fitbit_sleep_report",
]
candidate_read_only_tools = ["get_proactive_events", "get_sleep_context"]
candidate_read_only_tools = ["fitbit_health_snapshot", "fitbit_sleep_report"]
endpoint_env = [{env = "FITBIT_MONITOR_PORT", process = "monitor"}]
candidate_env = {FITBIT_BACKEND = "recording"}
59 changes: 59 additions & 0 deletions docs/content-context-v3.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
# Fitbit Content 与睡眠上下文

Fitbit v3.1 不再登记 proactive source,也不再用 MCP 工具搬运主动事件。插件只组合已有能力:

```text
Fitbit 公网
monitor(唯一采集 owner)
│ 本机 /api/agent,一次快照
FitbitContentRuntime
├── health events ──▶ Content.submit ──▶ Wake / Delivery
│ │
│ ▼
│ Content.unsettled ◀── delivered
│ │
│ ▼
│ monitor desired-state ACK
│ │
│ ▼
│ Content.ack
└── sleep ──▶ adapter.sqlite3 current cache
before Turn:仅 channel=wake 且未过期时追加 hint
```

## 不变量

1. `TIMERS` 只等待一次;插件在成功 tick 后持久化 `next_due` 并重新登记。
2. monitor 仍是 Fitbit 公网采集的唯一 owner;adapter 只读本机 `/api/agent`。
3. 健康事件先提交 Content,随后才原子更新插件私有的睡眠缓存与 `next_due`。
4. 睡眠不进入 Content,也不因状态变化单独唤醒;它只给现有 Wake Turn 增加上下文。
5. 外部 ACK 以“不再 pending”为成功事实。即使 ACK HTTP 返回后进程崩溃,下一轮也会确认事件已不在队列,再执行 `Content.ack`。
6. candidate 可以启动隔离 monitor 并完成 MCP handshake,但不会收到 `RUNTIME_STARTED`,因此不会登记 Timer、轮询 `/api/agent`、ACK 或写正式数据。

## 私有持久状态

`adapter.sqlite3/source_state` 只有一行:

- `next_due`:下一次本地 monitor 采集时间;
- `sleep_json`:最近一次睡眠判断;
- `sleep_observed_at` 与 `sleep_expires_at`:上下文新鲜度。

该文件不会复制 Content 的 item、delivery 或 ACK ledger。Content 继续独占这些权威事实。

## 保留与减少合同

| 对象 | owner | 正常增加 | 允许原位更新 | 逻辑失效 | 物理减少条件 | 恢复证据 |
|---|---|---|---|---|---|---|
| 已提交的健康事件、delivery 与 source ACK 完成事实 | Content | 新 item/revision 与 submission receipt 只追加 | 状态推进、selection、`settlement_ref` | `invalidated`、`abandoned`、`expired` 等 Content 状态 | 本插件没有物理减少协议 | Content row、原 payload、`settlement_ref` 与 `settled` 状态 |
| monitor pending / acked-id 队列投影 | monitor | 检测到事件时加入 pending;ACK 后加入 acked-id | monitor 现有检测状态与 pending 内容 | 事件过期或 ACK 后不再可投递 | ACK/过期可移除 pending;acked-id 超过既有固定上限可轮转 | 尚未提交前依赖 monitor 现有 state;提交后由 Content 成为唯一全量历史 owner |
| Session、Message 与 Turn | Core | 按 Core 协议追加 | Core 已批准的 metadata/terminal 状态 | 用户撤销或 Core 已定义的失效状态 | 只允许用户主动删除会话等 Core 已批准路径 | `sessions.db`、Message 与 Turn ledger |
| sleep current cache | Fitbit adapter | 首次建立 singleton | 每次成功 snapshot 覆盖 current payload 与时间 | `sleep_expires_at` 到期后不再注入 | 新 current snapshot 可以覆盖旧投影;没有历史裁切任务 | `sleep_observed_at`、`sleep_expires_at` 与当前 JSON |
| 纯诊断日志 | 各产生日志的 owner | 本实现不新增持久诊断日志 | 不适用 | 不适用 | 未来若新增,只能按固定数量轮转 | 固定轮转配置与当前日志文件 |

因此,“外部 ACK 成功”不会删除已发生的健康事实:monitor 的 pending 项可以消失,但 Content 的 settled row 仍保留原 payload 与 settlement receipt。adapter 不另造第二份历史或 ACK ledger。
75 changes: 75 additions & 0 deletions monitor/runtime_env.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,81 @@
from __future__ import annotations

import os
from pathlib import Path
from threading import Lock
from typing import TextIO


RUNTIME_LOG_MAX_BYTES = 1_048_576
RUNTIME_LOG_BACKUPS = 3


class RotatingTextLog:
"""Keep one diagnostic text log within fixed byte and generation limits."""

def __init__(
self,
path: Path,
*,
max_bytes: int = RUNTIME_LOG_MAX_BYTES,
backups: int = RUNTIME_LOG_BACKUPS,
) -> None:
if max_bytes <= 0 or backups < 0:
raise ValueError("runtime log rotation limits 无效")
self.path = path
self._max_bytes = max_bytes
self._backups = backups
self._lock = Lock()
path.parent.mkdir(parents=True, exist_ok=True)
self._stream = self._open()
self._size = path.stat().st_size if path.exists() else 0
if self._size >= self._max_bytes:
self._rotate()

def write(self, text: str) -> int:
input_length = len(text)
data = text.encode("utf-8")
with self._lock:
if self._size and self._size + len(data) > self._max_bytes:
self._rotate()
if len(data) > self._max_bytes:
data = data[-self._max_bytes :]
while data and (data[0] & 0xC0) == 0x80:
data = data[1:]
text = data.decode("utf-8")
self._stream.write(text)
self._size += len(data)
return input_length

def flush(self) -> None:
with self._lock:
self._stream.flush()

def close(self) -> None:
with self._lock:
self._stream.close()

def _rotate(self) -> None:
self._stream.close()
if self._backups > 0:
oldest = self.path.with_name(f"{self.path.name}.{self._backups}")
oldest.unlink(missing_ok=True)
for index in range(self._backups - 1, 0, -1):
source = self.path.with_name(f"{self.path.name}.{index}")
if source.exists():
os.replace(
source,
self.path.with_name(f"{self.path.name}.{index + 1}"),
)
if self.path.exists():
os.replace(self.path, self.path.with_name(f"{self.path.name}.1"))
else:
self.path.unlink(missing_ok=True)
self._stream = self._open()
self._size = 0

def _open(self) -> TextIO:
return self.path.open("a", encoding="utf-8", buffering=1)


def resolve_server_port(configured_port: int) -> int:
Expand Down
65 changes: 49 additions & 16 deletions monitor/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,13 +10,14 @@
from datetime import datetime, date, timedelta
from threading import Thread, Lock, Event
from pathlib import Path
from typing import TextIO
import tomllib
import requests as req
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from fastapi.responses import HTMLResponse, RedirectResponse, JSONResponse
import uvicorn

from runtime_env import resolve_server_port
from runtime_env import RotatingTextLog, resolve_server_port
import sleep_model
import retrain_guard
import build_sleep_diff_report
Expand Down Expand Up @@ -64,7 +65,7 @@ def writable(self):
return True


def _stream_points_to(path: Path, stream) -> bool:
def _stream_points_to(path: Path, stream: TextIO) -> bool:
try:
stream_fd = stream.fileno()
stream_stat = os.fstat(stream_fd)
Expand All @@ -77,26 +78,58 @@ def _stream_points_to(path: Path, stream) -> bool:
)


def _runtime_log_stream_fds(path: Path, streams: tuple[TextIO, ...]) -> set[int]:
"""Collect inherited descriptors that write directly to the runtime log."""

redirected_fds: set[int] = set()
for stream in streams:
if not _stream_points_to(path, stream):
continue
stream.flush()
redirected_fds.add(stream.fileno())
return redirected_fds


def _detach_stream_fds(stream_fds: set[int]) -> None:
"""Detach direct writers so all later text passes through the rotator."""

for stream_fd in stream_fds:
devnull_fd = os.open(os.devnull, os.O_WRONLY)
try:
os.dup2(devnull_fd, stream_fd)
finally:
os.close(devnull_fd)


def _install_runtime_log_mirror() -> None:
"""
保证无论通过何种方式启动,stdout/stderr 都会写入 monitor.runtime.log。
若上层已重定向到同一个文件,则不重复包裹,避免双写。
"""
if isinstance(sys.stdout, _TeeTextIO) or isinstance(sys.stderr, _TeeTextIO):
return
try:
RUNTIME_LOG_FILE.parent.mkdir(parents=True, exist_ok=True)
log_f = RUNTIME_LOG_FILE.open("a", encoding="utf-8", buffering=1)
except Exception:
"""Route stdout/stderr through the bounded runtime diagnostic log."""

# 1. 重复导入保持同一份 rotator;部分安装属于内部合同错误
installed = (
isinstance(sys.stdout, _TeeTextIO),
isinstance(sys.stderr, _TeeTextIO),
)
if installed == (True, True):
return
if any(installed):
raise RuntimeError("runtime log mirror 处于部分安装状态")

# 2. 先记住旧 inode 的直接写入者,再建立 rotator 并解除它们
RUNTIME_LOG_FILE.parent.mkdir(parents=True, exist_ok=True)
redirected_fds = _runtime_log_stream_fds(
RUNTIME_LOG_FILE,
(sys.stdout, sys.stderr),
)
log_f = RotatingTextLog(RUNTIME_LOG_FILE)
_detach_stream_fds(redirected_fds)
log_f.write(
f"\n===== fitbit-monitor start {datetime.now().strftime('%Y-%m-%d %H:%M:%S')} "
f"pid={os.getpid()} =====\n"
)
if not _stream_points_to(RUNTIME_LOG_FILE, sys.stdout):
sys.stdout = _TeeTextIO(sys.stdout, log_f)
if not _stream_points_to(RUNTIME_LOG_FILE, sys.stderr):
sys.stderr = _TeeTextIO(sys.stderr, log_f)

# 3. 两个流共享唯一 rotator;原始流继续承担终端或父进程可观察性
sys.stdout = _TeeTextIO(sys.stdout, log_f)
sys.stderr = _TeeTextIO(sys.stderr, log_f)
atexit.register(log_f.close)


Expand Down
Loading
Loading