Conversation
Implements Issue #102: background dispatcher polls schedules entering their reminder window, checks system calendar existence over WS before reminding when a calendar ref is bound, and cancels when the calendar was deleted. Adds per-device connection locking to fix a real concurrency gap between the new worker and the existing WS endpoint. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
There was a problem hiding this comment.
Review summary
Found correctness, delivery-reliability, authorization, and scaling issues in the new reminder dispatch path. See the inline comments for actionable details.
Verification: git diff --check passed. The focused pytest suite could not run because uv is unavailable in this environment; a direct reproduction confirmed the offset-aware/naive datetime failure.
| async def run_forever(self) -> None: | ||
| """周期性 tick,直到外部取消这个协程。""" | ||
| while True: | ||
| await self.tick(datetime.now()) |
There was a problem hiding this comment.
[P1] Use an aware clock here. Schedule timestamps are accepted and documented with offsets (for example +08:00), and datetime.fromisoformat() preserves that tzinfo. Passing naive datetime.now() into _time_window_reached raises TypeError: can't compare offset-naive and offset-aware datetimes; the exception escapes this loop and stops reminder dispatch entirely. Normalize both sides (preferably UTC) and add an offset-aware regression test.
| ref_id = schedule.system_schedule_ref_id | ||
| assert ref_id is not None | ||
| await self._connections.send( | ||
| schedule.user_id, |
There was a problem hiding this comment.
[P1] ConnectionManager is keyed by the WebSocket handshake device_id, but this lookup uses the schedule's user_id. There is no user-to-device registration or mapping in this PR, so a normal connection such as android_abc123 will not receive a schedule owned by default_user; send() simply returns False. The tests hide this by registering the user ID as though it were a device ID. Route through an actual device mapping (including multi-device behavior) or key connections consistently.
| snapshots = TimeWindowTriggerService(query_port).find_snapshots_entering_window(now) | ||
| for snapshot in snapshots: | ||
| try: | ||
| command_port.mark_triggered(snapshot.schedule_id, now) |
There was a problem hiding this comment.
[P1] This permanently stamps the row before any WebSocket delivery succeeds. If the target is offline, the connection key does not match, the socket send fails, or the process restarts while a refs check is pending, time_triggered_at remains set and the query excludes the schedule forever, so the reminder is lost. Only finalize the trigger after a durable/successful dispatch state, or persist a retryable outbox/state transition instead of committing first.
| if (result.ok and result.system_schedule_exists is not None) | ||
| else True | ||
| ) | ||
| await dispatcher.handle_refs_check_reply(result.schedule_id, calendar_exists) |
There was a problem hiding this comment.
[P1] The handler discards device_id and resolves pending state solely by client-supplied schedule_id. Any connected client that knows or guesses a pending ID can report system_schedule_exists=false and cancel another schedule. Bind the pending check to the replying device/user before popping it; the delete-ack handler below needs the same ownership/ref validation before clearing a row.
| return False | ||
| await connection.send_json(message) | ||
| async with self.lock_for(device_id): | ||
| await connection.send_json(message) |
There was a problem hiding this comment.
[P1] A stale/closed WebSocket can raise from send_json(), and this exception propagates through tick() into the unsupervised run_forever() task, permanently stopping all future reminders. This also contradicts the manager's best-effort contract. Re-read the current connection under the lock, handle send failures by unregistering only that stale instance, and keep the worker loop alive after transient failures.
|
|
||
| def list_schedules(self) -> Iterable[ScheduleSnapshot]: | ||
| """Return scheduled, not-yet-triggered rows with a non-null start_time.""" | ||
| statement = select(Schedule).where( |
There was a problem hiding this comment.
[P2] Every 30-second tick loads every future scheduled, untriggered time schedule into Python and only then checks whether its reminder window has arrived. This grows linearly with all outstanding schedules and the existing index cannot efficiently serve this predicate. Push a due-time bound into the query (or persist/index remind_at) so each poll only materializes candidates that can trigger now.
Summary
asyncio.LockinConnectionManagerto fix a real concurrency gap: the new background worker and the existing WS endpoint's own receive loop can now both write to the same connection, and Starlette'sWebSocket.send()has no internal locking.system_alarm_ref_id) is left unused this round.Test plan
ruff check .,mypy src/timeflow,pytestall pass (87 tests)uvicorn+ real WebSocket connections for all 8 acceptance criteria in Proposal:时间到达提醒完整闭环(时间监听 → 系统日历引用检查 → 提醒下发 + 日历清理) #102 except concurrent-write-to-same-connection (impractical to trigger deterministically by hand; covered bytest_send_serializes_concurrent_calls_to_the_same_deviceinstead)test_architecture.pyextended to assertinfrastructure/workersdoes not importdata/gateway/intelligence🤖 Generated with Claude Code