From cb4ea626ca586183db3495ec1fa62c3abda9e25a Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Sat, 4 Apr 2026 17:09:17 -0600 Subject: [PATCH] feat: add weekly schedule type for cogitate agents Weekly boundary-based scheduling with configurable day/time, dedup across supervisor restarts, and auto-registration of weekly-agents in config/schedules.json. Partner profile agent now runs weekly. --- AGENTS.md | 5 +- muse/partner.md | 2 +- tests/baselines/api/agents/agents-day.json | 4 +- tests/baselines/api/settings/providers.json | 4 +- tests/fixtures/journal/config/schedules.json | 6 + tests/test_generators.py | 2 +- tests/test_scheduler.py | 273 ++++++++++++++- think/dream.py | 340 +++++++++++++++++++ think/muse_cli.py | 3 +- think/scheduler.py | 156 ++++++++- 10 files changed, 771 insertions(+), 24 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 0f2ac5675..f6e61c241 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -289,15 +289,18 @@ Check and record onboarding state through the awareness system. Create facets an ## Identity Persistence -You maintain two files that give you continuity between sessions: +You maintain three files that give you continuity between sessions: - **`sol/self.md`** — Your identity file. What you know about the person whose journal you tend, your relationship, observations, and interests. Update when something genuinely changes your understanding. - **`sol/agency.md`** — Your initiative queue. Issues you've found, curation opportunities, follow-throughs. Update when you notice something worth tracking. +- **`sol/partner.md`** — Your understanding of the owner's behavioral patterns. Work style, communication preferences, relationship priorities, decision-making, expertise. Read-only in conversation — updated periodically by the partner profile agent. ### How to write Read current state: `sol call sol self` or `sol call sol agency` +Read partner profile: `sol call sol partner` (read-only — do not write in conversation) + Update a section of self.md (preferred — preserves other sections): ``` sol call sol self --update-section 'who I'\''m here for' --value 'Jer — founder-engineer, goes by Jer not Jeremie' diff --git a/muse/partner.md b/muse/partner.md index fba58e5c9..b84b4cc85 100644 --- a/muse/partner.md +++ b/muse/partner.md @@ -3,7 +3,7 @@ "title": "Partner Profile", "description": "Weekly observation of the journal owner's behavioral patterns — work style, communication, priorities, decision-making, expertise", - "schedule": "none", + "schedule": "weekly", "priority": 95, "instructions": {"system": "journal", "facets": true, "now": true} diff --git a/tests/baselines/api/agents/agents-day.json b/tests/baselines/api/agents/agents-day.json index 40eb0ecd8..27f27bc1e 100644 --- a/tests/baselines/api/agents/agents-day.json +++ b/tests/baselines/api/agents/agents-day.json @@ -292,7 +292,7 @@ "description": "Weekly observation of the journal owner's behavioral patterns \u2014 work style, communication, priorities, decision-making, expertise", "multi_facet": false, "output_format": null, - "schedule": "none", + "schedule": "weekly", "source": "system", "title": "Partner Profile", "type": "cogitate" @@ -476,4 +476,4 @@ } }, "runs": [] -} +} \ No newline at end of file diff --git a/tests/baselines/api/settings/providers.json b/tests/baselines/api/settings/providers.json index 59fe7e4b3..aba429cae 100644 --- a/tests/baselines/api/settings/providers.json +++ b/tests/baselines/api/settings/providers.json @@ -237,7 +237,7 @@ "disabled": false, "group": "Think", "label": "Partner Profile", - "schedule": "none", + "schedule": "weekly", "tier": 2, "type": "cogitate" }, @@ -469,4 +469,4 @@ "name": "openai" } ] -} +} \ No newline at end of file diff --git a/tests/fixtures/journal/config/schedules.json b/tests/fixtures/journal/config/schedules.json index a26d8fbae..e4789384c 100644 --- a/tests/fixtures/journal/config/schedules.json +++ b/tests/fixtures/journal/config/schedules.json @@ -1,5 +1,7 @@ { "daily_time": "03:00", + "weekly_day": "sunday", + "weekly_time": "03:00", "test:echo": { "cmd": ["sol", "echo", "-v"], "every": "hourly" @@ -8,6 +10,10 @@ "cmd": ["sol", "dream", "-v"], "every": "daily" }, + "test:weekly": { + "cmd": ["sol", "dream", "--weekly", "-v"], + "every": "weekly" + }, "test:disabled": { "cmd": ["sol", "noop"], "every": "hourly", diff --git a/tests/test_generators.py b/tests/test_generators.py index 6b13ea901..778f0495a 100644 --- a/tests/test_generators.py +++ b/tests/test_generators.py @@ -112,7 +112,7 @@ def test_scheduled_generators_have_valid_schedule(): muse = importlib.import_module("think.muse") generators = muse.get_muse_configs(type="generate") - valid_schedules = ("segment", "daily", "activity") + valid_schedules = ("segment", "daily", "activity", "weekly") for key, meta in generators.items(): sched = meta.get("schedule") diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index a693b56ba..6b5c4d488 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -52,6 +52,9 @@ def reset_scheduler_state(): mod._last_hour = None mod._daily_time = None mod._last_daily_mark = None + mod._weekly_day = None + mod._weekly_time = None + mod._last_weekly_mark = None yield mod._entries = {} mod._state = {} @@ -59,6 +62,9 @@ def reset_scheduler_state(): mod._last_hour = None mod._daily_time = None mod._last_daily_mark = None + mod._weekly_day = None + mod._weekly_time = None + mod._last_weekly_mark = None @pytest.fixture @@ -123,7 +129,7 @@ class TestLoadConfig: _write_config( journal_path, { - "bad": {"cmd": ["sol", "noop"], "every": "weekly"}, + "bad": {"cmd": ["sol", "noop"], "every": "biweekly"}, }, ) from think.scheduler import load_config @@ -435,6 +441,271 @@ class TestDailyTime: assert "daily_time" not in status[0] +class TestWeeklyTime: + """Tests for weekly scheduling — boundary computation and config parsing.""" + + def test_load_config_extracts_weekly_day_and_time(self, journal_path): + """load_config extracts weekly_day and weekly_time from schedules.json.""" + import think.scheduler as mod + + _write_config( + journal_path, + { + "weekly_day": "sunday", + "weekly_time": "04:00", + "w": {"cmd": ["sol", "dream", "--weekly"], "every": "weekly"}, + }, + ) + entries = mod.load_config() + assert "w" in entries + assert "weekly_day" not in entries + assert "weekly_time" not in entries + assert mod._weekly_day == "sunday" + assert mod._weekly_time == "04:00" + + def test_load_config_no_weekly_config(self, journal_path): + """When weekly_day/weekly_time are absent, globals are None.""" + import think.scheduler as mod + + _write_config(journal_path, {"a": {"cmd": ["sol", "x"], "every": "hourly"}}) + mod.load_config() + assert mod._weekly_day is None + assert mod._weekly_time is None + + def test_load_config_invalid_weekly_day(self, journal_path): + """Invalid weekly_day string is ignored.""" + import think.scheduler as mod + + _write_config( + journal_path, + { + "weekly_day": "notaday", + "w": {"cmd": ["sol", "x"], "every": "weekly"}, + }, + ) + mod.load_config() + assert mod._weekly_day is None + + def test_load_config_non_string_weekly_day(self, journal_path): + """Non-string weekly_day is ignored.""" + import think.scheduler as mod + + _write_config( + journal_path, + { + "weekly_day": 0, + "w": {"cmd": ["sol", "x"], "every": "weekly"}, + }, + ) + mod.load_config() + assert mod._weekly_day is None + + def test_weekly_day_case_insensitive(self, journal_path): + """Day name parsing is case-insensitive and accepts abbreviations.""" + import think.scheduler as mod + + for name in ["Sunday", "SUNDAY", "sun", "Sun"]: + _write_config( + journal_path, + { + "weekly_day": name, + "w": {"cmd": ["sol", "x"], "every": "weekly"}, + }, + ) + mod.load_config() + assert mod._weekly_day == name + + def test_compute_weekly_mark_past_boundary(self): + """When now is past this week's target, returns this week's boundary.""" + import think.scheduler as mod + + now = datetime(2026, 3, 22, 4, 0) + mark = mod._compute_weekly_mark(now, 6, "03:00") + assert mark == datetime(2026, 3, 22, 3, 0) + + def test_compute_weekly_mark_before_boundary(self): + """When now is before this week's target, returns last week's boundary.""" + import think.scheduler as mod + + now = datetime(2026, 3, 22, 2, 0) + mark = mod._compute_weekly_mark(now, 6, "03:00") + assert mark == datetime(2026, 3, 15, 3, 0) + + def test_compute_weekly_mark_midweek(self): + """Midweek, returns the most recent target day occurrence.""" + import think.scheduler as mod + + now = datetime(2026, 3, 25, 10, 0) + mark = mod._compute_weekly_mark(now, 6, "03:00") + assert mark == datetime(2026, 3, 22, 3, 0) + + def test_compute_weekly_mark_no_time_defaults_to_0300(self): + """When weekly_time is None, boundary defaults to 03:00.""" + import think.scheduler as mod + + now = datetime(2026, 3, 22, 4, 0) + mark = mod._compute_weekly_mark(now, 6, None) + assert mark == datetime(2026, 3, 22, 3, 0) + + def test_is_due_weekly_due(self): + """Weekly task is due when last_run is before the weekly boundary.""" + import think.scheduler as mod + + mod._weekly_day = "sunday" + mod._weekly_time = "03:00" + entry = {"cmd": ["sol", "x"], "every": "weekly"} + state = {"last_run": datetime(2026, 3, 21, 10, 0).timestamp()} + assert mod._is_due(entry, state, datetime(2026, 3, 22, 4, 0)) is True + + def test_is_due_weekly_not_due(self): + """Weekly task is not due when last_run is after the weekly boundary.""" + import think.scheduler as mod + + mod._weekly_day = "sunday" + mod._weekly_time = "03:00" + entry = {"cmd": ["sol", "x"], "every": "weekly"} + state = {"last_run": datetime(2026, 3, 22, 4, 0).timestamp()} + assert mod._is_due(entry, state, datetime(2026, 3, 25, 10, 0)) is False + + def test_is_due_weekly_no_state(self): + """Weekly task with no prior run is always due.""" + import think.scheduler as mod + + mod._weekly_day = "sunday" + entry = {"cmd": ["sol", "x"], "every": "weekly"} + assert mod._is_due(entry, None, datetime(2026, 3, 25, 10, 0)) is True + + def test_check_fires_at_weekly_boundary(self, journal_path): + """check() fires weekly tasks when the weekly boundary is crossed.""" + import think.scheduler as mod + + callosum = Mock() + callosum.emit = Mock(return_value=True) + + _write_config( + journal_path, + { + "weekly_day": "sunday", + "weekly_time": "03:00", + "w": {"cmd": ["sol", "dream", "--weekly"], "every": "weekly"}, + }, + ) + + mod.init(callosum) + + mod._last_hour = datetime(2026, 3, 21, 23, 0) + mod._last_daily_mark = datetime(2026, 3, 21, 0, 0) + mod._last_weekly_mark = datetime(2026, 3, 15, 3, 0) + + with _fake_now(datetime(2026, 3, 22, 3, 1)): + mod.check() + + callosum.emit.assert_called_once() + assert callosum.emit.call_args[1]["cmd"] == ["sol", "dream", "--weekly"] + + def test_check_no_fire_before_weekly_boundary(self, journal_path): + """check() does not fire weekly tasks before the weekly boundary.""" + import think.scheduler as mod + + callosum = Mock() + callosum.emit = Mock(return_value=True) + + _write_config( + journal_path, + { + "weekly_day": "sunday", + "weekly_time": "03:00", + "w": {"cmd": ["sol", "dream", "--weekly"], "every": "weekly"}, + }, + ) + + _write_state( + journal_path, + {"w": {"last_run": datetime(2026, 3, 15, 4, 0).timestamp()}}, + ) + + mod.init(callosum) + + mod._last_hour = datetime(2026, 3, 21, 22, 0) + mod._last_daily_mark = datetime(2026, 3, 21, 0, 0) + mod._last_weekly_mark = datetime(2026, 3, 15, 3, 0) + + with _fake_now(datetime(2026, 3, 21, 23, 1)): + mod.check() + + callosum.emit.assert_not_called() + + def test_missed_weeks_runs_once(self, journal_path): + """If supervisor was down for 3 weeks, weekly agent runs once on restart.""" + import think.scheduler as mod + + callosum = Mock() + callosum.emit = Mock(return_value=True) + + _write_config( + journal_path, + { + "weekly_day": "sunday", + "weekly_time": "03:00", + "w": {"cmd": ["sol", "dream", "--weekly"], "every": "weekly"}, + }, + ) + + _write_state( + journal_path, + {"w": {"last_run": datetime(2026, 3, 1, 4, 0).timestamp()}}, + ) + + mod.init(callosum) + + mod._last_hour = datetime(2026, 3, 22, 2, 0) + mod._last_daily_mark = datetime(2026, 3, 22, 0, 0) + mod._last_weekly_mark = datetime(2026, 3, 15, 3, 0) + + with _fake_now(datetime(2026, 3, 22, 3, 1)): + mod.check() + + callosum.emit.assert_called_once() + + def test_dedup_same_week_not_due(self): + """After running this week, weekly agent is not due again.""" + import think.scheduler as mod + + mod._weekly_day = "sunday" + mod._weekly_time = "03:00" + entry = {"cmd": ["sol", "x"], "every": "weekly"} + state = {"last_run": datetime(2026, 3, 22, 3, 30).timestamp()} + assert mod._is_due(entry, state, datetime(2026, 3, 26, 10, 0)) is False + + def test_format_next_due_weekly(self): + """_format_next_due shows next weekday and time.""" + import think.scheduler as mod + + mod._weekly_day = "sunday" + mod._weekly_time = "03:00" + entry = {"cmd": ["sol", "x"], "every": "weekly"} + state = {"last_run": datetime(2026, 3, 22, 4, 0).timestamp()} + now = datetime(2026, 3, 25, 10, 0) + + result = mod._format_next_due(entry, state, now) + assert "Sunday" in result + assert "03:00" in result + + def test_collect_status_includes_weekly_fields(self, journal_path): + """collect_status includes weekly_day and weekly_time for weekly entries.""" + import think.scheduler as mod + + mod._weekly_day = "sunday" + mod._weekly_time = "04:00" + mod._entries = {"w": {"cmd": ["sol", "x"], "every": "weekly"}} + mod._state = {} + + status = mod.collect_status() + assert len(status) == 1 + assert status[0]["weekly_day"] == "sunday" + assert status[0]["weekly_time"] == "04:00" + + # --------------------------------------------------------------------------- # init # --------------------------------------------------------------------------- diff --git a/think/dream.py b/think/dream.py index 931cd78fd..49b06c902 100644 --- a/think/dream.py +++ b/think/dream.py @@ -1031,6 +1031,302 @@ def run_daily_prompts( return (total_success, total_failed, all_failed_names) +def run_weekly_prompts( + day: str, + refresh: bool, + verbose: bool, + max_concurrency: int = 2, + stream: str | None = None, + timeout: int | None = 610, +) -> tuple[int, int, list[str]]: + """Run all weekly scheduled prompts in priority order. + + Loads all weekly prompts, groups by priority, and executes each group with + bounded concurrency. Structurally identical to run_daily_prompts but for + weekly-scheduled agents (e.g., partner profile). + + Args: + day: Day in YYYYMMDD format (reference day for agent context) + refresh: Whether to regenerate existing outputs + verbose: Verbose logging + max_concurrency: Max agents to run concurrently per priority group. + 0 means unlimited (all agents in a group run in parallel). + + Returns: + Tuple of (success_count, fail_count, failed_names). + """ + target_schedule = "weekly" + + # Load ALL scheduled prompts (both generators and agents) + all_prompts = get_muse_configs(schedule=target_schedule) + + if not all_prompts: + logging.info(f"No prompts found for schedule: {target_schedule}") + return (0, 0, []) + + # Group prompts by priority + priority_groups: dict[int, list[tuple[str, dict]]] = {} + for name, config in all_prompts.items(): + priority = config["priority"] # Required field, validated by get_muse_configs + priority_groups.setdefault(priority, []).append((name, config)) + + # Pre-compute shared data for multi-facet prompts + day_formatted = iso_date(day) + input_summary = day_input_summary(day) + enabled_facets = get_enabled_facets() + active_facets = get_active_facets(day) + + total_prompts = sum(len(prompts) for prompts in priority_groups.values()) + num_groups = len(priority_groups) + _update_status( + mode=target_schedule, + day=day, + stream=stream, + agents_total=total_prompts, + agents_completed=0, + current_agents=[], + ) + + logging.info( + f"Running {total_prompts} prompts for {day} in {num_groups} priority groups" + ) + + emit( + "started", + mode=target_schedule, + day=day, + count=total_prompts, + groups=num_groups, + ) + + start_time = time.time() + total_success = 0 + total_failed = 0 + all_failed_names: list[str] = [] + + # Process each priority group in order + for priority in sorted(priority_groups.keys()): + prompts_list = priority_groups[priority] + _update_status(current_group_priority=priority) + logging.info(f"Starting priority {priority} ({len(prompts_list)} prompts)") + + emit( + "group_started", + mode=target_schedule, + day=day, + priority=priority, + count=len(prompts_list), + ) + + spawned: list[ + tuple[str, str, dict, str | None] + ] = [] # (agent_id, name, config, facet) + group_success = 0 + group_failed = 0 + + for prompt_name, config in prompts_list: + is_generate = config["type"] == "generate" + + # Check exclude_streams filter + exclude_patterns = config.get("exclude_streams") + if exclude_patterns and stream: + if any(fnmatch.fnmatch(stream, pat) for pat in exclude_patterns): + logging.info( + f"Skipping {prompt_name}: stream '{stream}' matches exclude_streams" + ) + continue + + try: + if config.get("multi_facet"): + always_run = config.get("always", False) + + for facet_name in enabled_facets.keys(): + if not always_run and facet_name not in active_facets: + logging.info( + f"Skipping {prompt_name} for {facet_name}: " + f"no activity on {day_formatted}" + ) + continue + + logging.info(f"Spawning {prompt_name} for facet: {facet_name}") + + # Always pass day for instructions.day context + request_config: dict = {"facet": facet_name, "day": day} + if is_generate: + request_config["output"] = config.get("output", "md") + if refresh: + request_config["refresh"] = True + elif config.get("output"): + # Cogitate agents with explicit output get auto-persisted + request_config["output"] = config["output"] + env: dict[str, str] = { + "SOL_DAY": day, + "SOL_FACET": facet_name, + } + request_config["env"] = env + request_config["schedule"] = target_schedule + + prompt = ( + "" + if is_generate + else f"Processing facet '{facet_name}' for {day_formatted}: {input_summary}. Use get_facet('{facet_name}') to load context." + ) + + agent_id = _cortex_request_with_retry( + prompt=prompt, + name=prompt_name, + config=request_config, + ) + if agent_id is None: + group_failed += 1 + all_failed_names.append( + f"{prompt_name}/{facet_name} (send)" + ) + continue + spawned.append((agent_id, prompt_name, config, facet_name)) + emit( + "agent_started", + mode=target_schedule, + day=day, + name=prompt_name, + agent_id=agent_id, + facet=facet_name, + ) + logging.info( + f"Started {prompt_name} for {facet_name} (ID: {agent_id})" + ) + + # Drain batch when concurrency limit reached + if max_concurrency and len(spawned) >= max_concurrency: + _update_status( + current_agents=[name for _, name, _, _ in spawned] + ) + s, f, fn = _drain_priority_batch( + spawned, + target_schedule, + day, + None, + stream, + timeout, + ) + group_success += s + group_failed += f + all_failed_names.extend(fn) + spawned = [] + _update_status( + agents_completed=total_success + + total_failed + + group_success + + group_failed, + current_agents=[], + ) + else: + # Regular single-instance prompt + logging.info(f"Spawning {prompt_name}") + + # Always pass day for instructions.day context + request_config: dict = {"day": day} + if is_generate: + request_config["output"] = config.get("output", "md") + if refresh: + request_config["refresh"] = True + env: dict[str, str] = {"SOL_DAY": day} + request_config["env"] = env + request_config["schedule"] = target_schedule + + prompt = ( + "" + if is_generate + else f"Running scheduled task for {day_formatted}: {input_summary}." + ) + + agent_id = _cortex_request_with_retry( + prompt=prompt, + name=prompt_name, + config=request_config, + ) + if agent_id is None: + group_failed += 1 + all_failed_names.append(f"{prompt_name} (send)") + continue + spawned.append((agent_id, prompt_name, config, None)) + emit( + "agent_started", + mode=target_schedule, + day=day, + name=prompt_name, + agent_id=agent_id, + ) + logging.info(f"Started {prompt_name} (ID: {agent_id})") + + # Drain batch when concurrency limit reached + if max_concurrency and len(spawned) >= max_concurrency: + _update_status( + current_agents=[name for _, name, _, _ in spawned] + ) + s, f, fn = _drain_priority_batch( + spawned, target_schedule, day, None, stream, timeout + ) + group_success += s + group_failed += f + all_failed_names.extend(fn) + spawned = [] + _update_status( + agents_completed=total_success + + total_failed + + group_success + + group_failed, + current_agents=[], + ) + + except Exception as e: + logging.error(f"Failed to spawn {prompt_name}: {e}") + group_failed += 1 + all_failed_names.append(f"{prompt_name} (spawn)") + + # Drain any remaining agents in this priority group + _update_status(current_agents=[name for _, name, _, _ in spawned]) + s, f, fn = _drain_priority_batch( + spawned, target_schedule, day, None, stream, timeout + ) + group_success += s + group_failed += f + all_failed_names.extend(fn) + _update_status( + agents_completed=total_success + + total_failed + + group_success + + group_failed, + current_agents=[], + ) + + total_success += group_success + total_failed += group_failed + + emit( + "group_completed", + mode=target_schedule, + day=day, + priority=priority, + success=group_success, + failed=group_failed, + ) + + duration_ms = int((time.time() - start_time) * 1000) + emit( + "completed", + mode=target_schedule, + day=day, + success=total_success, + failed=total_failed, + failed_names=all_failed_names, + duration_ms=duration_ms, + ) + + logging.info(f"Prompts completed: {total_success} succeeded, {total_failed} failed") + return (total_success, total_failed, all_failed_names) + + def run_activity_prompts( day: str, activity_id: str, @@ -1534,6 +1830,7 @@ def dry_run( flush: bool = False, refresh: bool = False, stream: str | None = None, + weekly: bool = False, ) -> None: """Print what dream would execute without spawning any agents.""" day_formatted = iso_date(day) @@ -1590,6 +1887,15 @@ def dry_run( _dry_run_flush(day, segment or "") return + if weekly: + all_prompts = get_muse_configs(schedule="weekly") + print(f"Day {day_formatted} — weekly agents\n") + if not all_prompts: + print("No prompts for schedule: weekly") + else: + _print_prompt_table(all_prompts, day, refresh=refresh, stream=stream) + return + if segments: segs = cluster_segments(day) if not segs: @@ -1867,6 +2173,11 @@ def parse_args() -> argparse.ArgumentParser: action="store_true", help="List days with pending daily processing and exit", ) + parser.add_argument( + "--weekly", + action="store_true", + help="Run weekly-scheduled agents (incompatible with --segment, --segments, --activity, --flush)", + ) parser.add_argument( "--dry-run", action="store_true", @@ -1949,6 +2260,11 @@ def main() -> None: if args.segments and (args.segment or args.facet): parser.error("--segments is incompatible with --segment and --facet") + if args.weekly and (args.segment or args.segments or args.activity or args.flush): + parser.error( + "--weekly is incompatible with --segment, --segments, --activity, and --flush" + ) + if args.dry_run: dry_run( day, @@ -1959,6 +2275,7 @@ def main() -> None: flush=args.flush, refresh=args.refresh, stream=args.stream, + weekly=args.weekly, ) sys.exit(0) @@ -2075,6 +2392,29 @@ def main() -> None: start_time = time.time() + # Handle weekly mode — dispatch weekly agents, no pre/post phases + if args.weekly: + success_count, fail_count, failed_names = run_weekly_prompts( + day=day, + refresh=args.refresh, + verbose=args.verbose, + max_concurrency=args.jobs, + stream=args.stream, + ) + + duration_ms = int((time.time() - start_time) * 1000) + logging.info( + f"Weekly dream completed in {duration_ms}ms: " + f"{success_count} succeeded, {fail_count} failed" + ) + day_log(day, f"dream --weekly failed={fail_count}") + + if fail_count > 0: + names = ", ".join(failed_names) + logging.error(f"{fail_count} weekly prompt(s) failed: {names}") + sys.exit(1) + sys.exit(0) + # PRE-PHASE: Run sense repair (daily only) if not args.segment: logging.info("Running pre-phase: sense repair") diff --git a/think/muse_cli.py b/think/muse_cli.py index c375be83d..eaf1abb8b 100644 --- a/think/muse_cli.py +++ b/think/muse_cli.py @@ -207,13 +207,14 @@ def list_prompts( groups: dict[str, list[tuple[str, dict[str, Any]]]] = { "segment": [], "daily": [], + "weekly": [], "activity": [], "unscheduled": [], } for key, info in sorted(configs.items()): sched = info.get("schedule") - if sched in ("segment", "daily", "activity"): + if sched in ("segment", "daily", "weekly", "activity"): groups[sched].append((key, info)) else: groups["unscheduled"].append((key, info)) diff --git a/think/scheduler.py b/think/scheduler.py index 4e5281092..cf4ec48e9 100644 --- a/think/scheduler.py +++ b/think/scheduler.py @@ -27,7 +27,7 @@ from think.utils import get_journal, now_ms, setup_cli logger = logging.getLogger(__name__) # Valid schedule intervals -INTERVALS = {"hourly", "daily"} +INTERVALS = {"hourly", "daily", "weekly"} # --------------------------------------------------------------------------- # Module state (populated by init(), used by check()) @@ -38,6 +38,9 @@ _callosum: Any = None # CallosumConnection _last_hour: datetime | None = None _daily_time: str | None = None _last_daily_mark: datetime | None = None +_weekly_day: str | None = None +_weekly_time: str | None = None +_last_weekly_mark: datetime | None = None # --------------------------------------------------------------------------- @@ -47,11 +50,13 @@ _last_daily_mark: datetime | None = None def load_config() -> dict[str, dict[str, Any]]: """Read config/schedules.json and return validated entries.""" - global _daily_time + global _daily_time, _weekly_day, _weekly_time config_path = Path(get_journal()) / "config" / "schedules.json" if not config_path.exists(): _daily_time = None + _weekly_day = None + _weekly_time = None return {} try: @@ -60,6 +65,8 @@ def load_config() -> dict[str, dict[str, Any]]: except (json.JSONDecodeError, OSError) as exc: logger.warning("Failed to load schedules config: %s", exc) _daily_time = None + _weekly_day = None + _weekly_time = None return {} if not isinstance(raw, dict): @@ -67,6 +74,8 @@ def load_config() -> dict[str, dict[str, Any]]: "schedules.json must be a JSON object, got %s", type(raw).__name__ ) _daily_time = None + _weekly_day = None + _weekly_time = None return {} # Extract daily_time metadata (not a schedule entry) @@ -75,6 +84,23 @@ def load_config() -> dict[str, dict[str, Any]]: logger.warning("schedules.json: daily_time must be a string, ignoring") _daily_time = None + # Extract weekly_day metadata + _weekly_day = raw.pop("weekly_day", None) + if _weekly_day is not None and not isinstance(_weekly_day, str): + logger.warning("schedules.json: weekly_day must be a string, ignoring") + _weekly_day = None + elif _weekly_day is not None and _parse_weekly_day(_weekly_day) is None: + logger.warning( + "schedules.json: unrecognized weekly_day '%s', ignoring", _weekly_day + ) + _weekly_day = None + + # Extract weekly_time metadata + _weekly_time = raw.pop("weekly_time", None) + if _weekly_time is not None and not isinstance(_weekly_time, str): + logger.warning("schedules.json: weekly_time must be a string, ignoring") + _weekly_time = None + entries: dict[str, dict[str, Any]] = {} for name, entry in raw.items(): if not isinstance(entry, dict): @@ -161,6 +187,31 @@ def _parse_daily_time(raw: str | None) -> tuple[int, int] | None: return None +DAY_NAMES: dict[str, int] = { + "monday": 0, + "mon": 0, + "tuesday": 1, + "tue": 1, + "wednesday": 2, + "wed": 2, + "thursday": 3, + "thu": 3, + "friday": 4, + "fri": 4, + "saturday": 5, + "sat": 5, + "sunday": 6, + "sun": 6, +} + + +def _parse_weekly_day(raw: str | None) -> int | None: + """Parse day-of-week name. Returns weekday int (0=Monday, 6=Sunday) or None.""" + if not raw or not isinstance(raw, str): + return None + return DAY_NAMES.get(raw.strip().lower()) + + def _compute_daily_mark(now: datetime, daily_time_str: str | None) -> datetime: """Compute the most recent daily boundary datetime. @@ -178,6 +229,27 @@ def _compute_daily_mark(now: datetime, daily_time_str: str | None) -> datetime: return today_mark - timedelta(days=1) +def _compute_weekly_mark( + now: datetime, weekly_day: int, weekly_time_str: str | None +) -> datetime: + """Compute the most recent weekly boundary datetime. + + Returns the most recent occurrence of the target weekday at the target time. + If now is past this week's boundary, returns this week's. Otherwise last week's. + """ + parsed = _parse_daily_time(weekly_time_str) + if parsed is None: + h, m = 3, 0 # default 03:00 + else: + h, m = parsed + days_since = (now.weekday() - weekly_day) % 7 + target_date = now - timedelta(days=days_since) + target_mark = target_date.replace(hour=h, minute=m, second=0, microsecond=0) + if now >= target_mark: + return target_mark + return target_mark - timedelta(weeks=1) + + def _is_due(entry: dict, state_entry: dict | None, now: datetime) -> bool: """Check if an entry is due based on its interval and last_run.""" last_run = (state_entry or {}).get("last_run") @@ -194,6 +266,11 @@ def _is_due(entry: dict, state_entry: dict | None, now: datetime) -> bool: return last_dt < _hour_mark(now) if every == "daily": return last_dt < _compute_daily_mark(now, _daily_time) + if every == "weekly": + weekly_day_val = _parse_weekly_day(_weekly_day) + if weekly_day_val is None: + weekly_day_val = 6 # default Sunday + return last_dt < _compute_weekly_mark(now, weekly_day_val, _weekly_time) return False @@ -204,7 +281,7 @@ def _is_due(entry: dict, state_entry: dict | None, now: datetime) -> bool: def init(callosum: Any) -> None: """Initialize scheduler with a Callosum connection. Load config and state.""" - global _entries, _state, _callosum, _last_hour, _last_daily_mark + global _entries, _state, _callosum, _last_hour, _last_daily_mark, _last_weekly_mark _callosum = callosum _entries = load_config() @@ -213,6 +290,10 @@ def init(callosum: Any) -> None: now = datetime.now() _last_hour = _hour_mark(now) _last_daily_mark = _compute_daily_mark(now, _daily_time) + weekly_day_val = _parse_weekly_day(_weekly_day) + if weekly_day_val is None: + weekly_day_val = 6 + _last_weekly_mark = _compute_weekly_mark(now, weekly_day_val, _weekly_time) if _entries: logger.info( @@ -232,7 +313,10 @@ def register_defaults() -> None: """ global _entries - if "heartbeat" in _entries: + need_heartbeat = "heartbeat" not in _entries + need_weekly = "weekly-agents" not in _entries + + if not need_heartbeat and not need_weekly: return # Read raw config (preserving daily_time and other entries) @@ -251,14 +335,26 @@ def register_defaults() -> None: if not isinstance(raw, dict): raw = {} - if "heartbeat" in raw: - return # Already in config file — don't overwrite user customization + changed = False + + if need_heartbeat and "heartbeat" not in raw: + raw["heartbeat"] = { + "cmd": ["sol", "heartbeat"], + "every": "daily", + "enabled": True, + } + changed = True - raw["heartbeat"] = { - "cmd": ["sol", "heartbeat"], - "every": "daily", - "enabled": True, - } + if need_weekly and "weekly-agents" not in raw: + raw["weekly-agents"] = { + "cmd": ["sol", "dream", "--weekly", "-v"], + "every": "weekly", + "enabled": True, + } + changed = True + + if not changed: + return # Atomic write fd, tmp_path = tempfile.mkstemp(dir=config_dir, suffix=".tmp", prefix=".schedules_") @@ -267,7 +363,7 @@ def register_defaults() -> None: with open(fd, "w", encoding="utf-8") as f: json.dump(raw, f, indent=2) tmp_file.replace(config_path) - logger.info("Auto-registered heartbeat schedule in config/schedules.json") + logger.info("Auto-registered default schedule(s) in config/schedules.json") except BaseException: tmp_file.unlink(missing_ok=True) raise @@ -282,7 +378,7 @@ def check() -> None: Called each supervisor tick (~1s). Does nothing unless an hour or day boundary has been crossed since the last check. """ - global _entries, _last_hour, _last_daily_mark + global _entries, _last_hour, _last_daily_mark, _last_weekly_mark if _last_hour is None: return @@ -290,11 +386,16 @@ def check() -> None: now = datetime.now() current_hour = _hour_mark(now) current_daily_mark = _compute_daily_mark(now, _daily_time) + weekly_day_val = _parse_weekly_day(_weekly_day) + if weekly_day_val is None: + weekly_day_val = 6 + current_weekly_mark = _compute_weekly_mark(now, weekly_day_val, _weekly_time) hour_changed = current_hour != _last_hour daily_mark_changed = current_daily_mark != _last_daily_mark + weekly_mark_changed = current_weekly_mark != _last_weekly_mark - if not hour_changed and not daily_mark_changed: + if not hour_changed and not daily_mark_changed and not weekly_mark_changed: return # Boundary crossed — reload config for freshest definitions @@ -305,6 +406,13 @@ def check() -> None: if new_daily_mark != _last_daily_mark: daily_mark_changed = True _last_daily_mark = new_daily_mark + new_weekly_day_val = _parse_weekly_day(_weekly_day) + if new_weekly_day_val is None: + new_weekly_day_val = 6 + new_weekly_mark = _compute_weekly_mark(now, new_weekly_day_val, _weekly_time) + if new_weekly_mark != _last_weekly_mark: + weekly_mark_changed = True + _last_weekly_mark = new_weekly_mark if not _entries: return @@ -318,6 +426,8 @@ def check() -> None: continue if every == "daily" and not daily_mark_changed: continue + if every == "weekly" and not weekly_mark_changed: + continue if not _is_due(entry, _state.get(name), now): continue @@ -365,6 +475,11 @@ def collect_status() -> list[dict[str, Any]]: } if entry["every"] == "daily" and _daily_time: entry_status["daily_time"] = _daily_time + if entry["every"] == "weekly": + if _weekly_day: + entry_status["weekly_day"] = _weekly_day + if _weekly_time: + entry_status["weekly_time"] = _weekly_time result.append(entry_status) return result @@ -396,6 +511,13 @@ def _format_next_due(entry: dict, state_entry: dict | None, now: datetime) -> st if every == "daily": parsed = _parse_daily_time(_daily_time) return f"{parsed[0]:02d}:{parsed[1]:02d}" if parsed else "midnight" + if every == "weekly": + weekly_day_val = _parse_weekly_day(_weekly_day) + if weekly_day_val is None: + weekly_day_val = 6 + weekly_mark = _compute_weekly_mark(now, weekly_day_val, _weekly_time) + nxt = weekly_mark + timedelta(weeks=1) + return f"{nxt.strftime('%A')} {nxt.strftime('%H:%M')}" return "?" @@ -419,9 +541,13 @@ def main() -> None: return # Extract daily_time metadata before processing entries - global _daily_time + global _daily_time, _weekly_day, _weekly_time raw_daily_time = config.pop("daily_time", None) _daily_time = raw_daily_time if isinstance(raw_daily_time, str) else None + raw_weekly_day = config.pop("weekly_day", None) + raw_weekly_time = config.pop("weekly_time", None) + _weekly_day = raw_weekly_day if isinstance(raw_weekly_day, str) else None + _weekly_time = raw_weekly_time if isinstance(raw_weekly_time, str) else None if not config: print("No schedules configured.") -- 2.51.2