@@ -509,3 +509,158 @@ def test_swept_conversation_can_start_a_new_run(self, user_id) -> None:
509509 finally :
510510 db .close ()
511511 _wait_status (controller , run .id , {"done" })
512+
513+
514+ class TestInstanceScopedReconciliation :
515+ """A restart must fail only THIS instance's orphaned runs, never a
516+ different live instance's healthy runs."""
517+
518+ def test_owner_instance_is_stamped_while_running (self , user_id ) -> None :
519+ from app .models import Run
520+
521+ uid , cid = user_id
522+ controller = Controller (runner = StubRunner (stdout = RESULT_OK ), instance_id = "inst-A" )
523+ db = get_session_factory ()()
524+ try :
525+ run = controller .start_run (db , user_id = uid , conversation_id = cid , prompt = "hi" )
526+ run_id = run .id
527+ finally :
528+ db .close ()
529+ _wait_status (controller , run_id , {"done" })
530+ # finished run releases ownership
531+ db = get_session_factory ()()
532+ try :
533+ assert db .get (Run , run_id ).owner_instance is None
534+ finally :
535+ db .close ()
536+
537+ def test_reconcile_ignores_other_instances_runs (self , user_id ) -> None :
538+ from app .models import Run
539+
540+ _uid , cid = user_id
541+ db = get_session_factory ()()
542+ try :
543+ # a run owned by a DIFFERENT, still-live instance
544+ db .add (
545+ Run (
546+ id = "run_other" ,
547+ conversation_id = cid ,
548+ prompt = "p" ,
549+ status = "running" ,
550+ owner_instance = "inst-B" ,
551+ )
552+ )
553+ db .commit ()
554+ finally :
555+ db .close ()
556+
557+ # instance A restarts and reconciles: must NOT touch inst-B's run
558+ n = Controller (runner = StubRunner (), instance_id = "inst-A" ).reconcile_orphaned_runs ()
559+ assert n == 0
560+
561+ db = get_session_factory ()()
562+ try :
563+ run = db .get (Run , "run_other" )
564+ assert run .status == "running" # left alone
565+ assert run .owner_instance == "inst-B"
566+ finally :
567+ db .close ()
568+
569+ def test_reconcile_claims_own_and_null_owner_runs (self , user_id ) -> None :
570+ from app .models import Run
571+
572+ _uid , cid = user_id
573+ # a second conversation so two running rows can coexist (the
574+ # partial unique index is per-conversation)
575+ db = get_session_factory ()()
576+ try :
577+ other_conv = new_id ("cnv" )
578+ db .add (Conversation (id = other_conv , user_id = _uid , title = "t2" ))
579+ db .add (
580+ Run (
581+ id = "run_mine" ,
582+ conversation_id = cid ,
583+ prompt = "p" ,
584+ status = "running" ,
585+ owner_instance = "inst-A" ,
586+ )
587+ )
588+ db .add (
589+ Run (
590+ id = "run_legacy" ,
591+ conversation_id = other_conv ,
592+ prompt = "p" ,
593+ status = "running" ,
594+ owner_instance = None , # pre-migration / unclaimed
595+ )
596+ )
597+ db .commit ()
598+ finally :
599+ db .close ()
600+
601+ n = Controller (runner = StubRunner (), instance_id = "inst-A" ).reconcile_orphaned_runs ()
602+ assert n == 2 # own run + the NULL-owner legacy run
603+
604+ db = get_session_factory ()()
605+ try :
606+ assert db .get (Run , "run_mine" ).status == "error"
607+ assert db .get (Run , "run_legacy" ).status == "error"
608+ # ownership released on reconcile
609+ assert db .get (Run , "run_mine" ).owner_instance is None
610+ finally :
611+ db .close ()
612+
613+
614+ class TestDrain :
615+ """Graceful shutdown: reject new runs and wait for in-flight ones."""
616+
617+ def test_drain_waits_for_in_flight_run (self , user_id ) -> None :
618+ from app .models import Run
619+
620+ uid , cid = user_id
621+ # a run that takes a moment; drain must block until it finishes
622+ controller = Controller (runner = StubRunner (delay = 0.3 , stdout = RESULT_OK ))
623+ db = get_session_factory ()()
624+ try :
625+ run = controller .start_run (db , user_id = uid , conversation_id = cid , prompt = "hi" )
626+ run_id = run .id
627+ finally :
628+ db .close ()
629+
630+ still = controller .drain (timeout = 10 )
631+ assert still == 0 # finished within the window
632+
633+ db = get_session_factory ()()
634+ try :
635+ assert db .get (Run , run_id ).status == "done"
636+ finally :
637+ db .close ()
638+
639+ def test_drain_rejects_new_runs (self , user_id ) -> None :
640+ uid , cid = user_id
641+ controller = Controller (runner = StubRunner (stdout = RESULT_OK ))
642+ controller .drain (timeout = 1 ) # sets draining
643+ db = get_session_factory ()()
644+ try :
645+ with pytest .raises (RuntimeError , match = "shutting down" ):
646+ controller .start_run (db , user_id = uid , conversation_id = cid , prompt = "late" )
647+ finally :
648+ db .close ()
649+
650+ def test_drain_returns_count_still_running_at_deadline (self , user_id ) -> None :
651+ uid , cid = user_id
652+ # a run longer than the drain window: still in flight at deadline
653+ controller = Controller (runner = StubRunner (delay = 2.0 , stdout = RESULT_OK ))
654+ db = get_session_factory ()()
655+ try :
656+ run = controller .start_run (db , user_id = uid , conversation_id = cid , prompt = "slow" )
657+ run_id = run .id
658+ finally :
659+ db .close ()
660+
661+ still = controller .drain (timeout = 0.2 )
662+ assert still == 1 # did not finish in time
663+
664+ # it is still owned by this instance and would be reconciled on
665+ # the next startup
666+ _wait_status (controller , run_id , {"done" }) # let it finish to clean up
0 commit comments